- 745
Some((entry.id.clone(), earlier.digest())) - 746
} - 747
_ => None, - 748
}); - 749
let parent = self.tail_id.clone(); - 750
let payload = match previous { - 751
Some((entry, earlier)) if earlier == digest => { - 752
EntryPayload::TurnCapabilitiesRef(crate::types::TurnCapabilitiesRef { - 753
entry, - 754
digest, - 755
epoch: bound.epoch, - 756
}) - 757
} - 758
_ => EntryPayload::TurnCapabilitiesBound(bound), - 759
}; - 760
self.append(Entry::new(parent, payload)) - 761
} - 762
- 763
pub fn append_work(&mut self, event: WorkEvent) -> Result<Entry, SessionError> { - 764
let parent = self.tail_id.clone(); - 765
let candidate = Entry::new(parent, EntryPayload::Work(event)); - 766
let mut chain = self.chain_to_root(); - 767
chain.push(&candidate); - 768
crate::work::project_work(&chain).map_err(|error| SessionError::Corrupt { - 769
line: 0, - 770
message: format!("invalid work event: {error}"), - 771
})?; - 772
let appended = self.append(candidate)?; - 773
self.promote_ready_work_items()?; - 774
Ok(appended) - 775
} - 776
- 777
fn promote_ready_work_items(&mut self) -> Result<(), SessionError> { - 778
let Some(projection) = self - 779
.work_projection() - 780
.map_err(|error| SessionError::Corrupt { - 781
line: 0, - 782
message: error.to_string(), - 783
})? - 784
else { - 785
return Ok(()); - 786
}; - 787
if projection.status != crate::types::WorkContractStatus::Active { - 788
return Ok(()); - 789
} - 790
let ready: Vec<String> = projection - 791
.contract - 792
.items - 793
.iter() - 794
.filter(|definition| { - 795
projection - 796
.items - 797
.get(&definition.item_id) - 798
.is_some_and(|state| state.status == crate::types::WorkItemStatus::Proposed) - 799
&& definition.dependencies.iter().all(|dependency| { - 800
projection.items.get(dependency).is_some_and(|state| { - 801
matches!( - 802
state.status, - 803
crate::types::WorkItemStatus::Succeeded - 804
| crate::types::WorkItemStatus::Skipped - 805
) - 806
}) - 807
}) - 808
}) - 809
.map(|definition| definition.item_id.clone()) - 810
.collect(); - 811
for item_id in ready { - 812
let current = self - 813
.work_projection() - 814
.map_err(|error| SessionError::Corrupt { - 815
line: 0, - 816
message: error.to_string(), - 817
})?; - 818
let Some(current) = current else { break }; - 819
self.append_work(crate::types::WorkEvent { - 820
contract_id: current.contract.contract_id, - 821
revision: current.contract.revision, - 822
kind: crate::types::WorkEventKind::ItemStatusChanged { - 823
item_id, - 824
from: crate::types::WorkItemStatus::Proposed, - 825
to: crate::types::WorkItemStatus::Ready, - 826
attempt: 0, - 827
reason: "dependencies satisfied".into(), - 828
}, - 829
})?; - 830
} - 831
Ok(()) - 832
} - 833
- 834
pub fn work_projection( - 835
&self, - 836
) -> Result<Option<crate::work::WorkProjection>, crate::work::WorkError> { - 837
crate::work::project_work(&self.chain_to_root()) - 838
} - 839
- 840
pub fn evidence_exists(&self, evidence: &crate::types::EvidenceRef) -> bool { - 841
let session_id = self.header().map(|header| header.session_id.as_str()); - 842
self.chain_to_root().iter().any(|entry| match evidence { - 843
crate::types::EvidenceRef::LedgerEntry { - 844
session_id: evidence_session, - 845
entry_id, - 846
} - 847
| crate::types::EvidenceRef::Receipt { - 848
session_id: evidence_session, - 849
entry_id, - 850
} => session_id == Some(evidence_session.as_str()) && entry.id == *entry_id, - 851
crate::types::EvidenceRef::ToolResult { - 852
session_id: evidence_session, - 853
tool_use_id, - 854
} => { - 855
session_id == Some(evidence_session.as_str()) - 856
&& matches!( - 857
&entry.payload, - 858
crate::types::EntryPayload::Message(record) - 859
if record.message.content.iter().any(|block| matches!( - 860
block, - 861
vak_llm::ContentBlock::ToolResult { tool_use_id: id, .. } - 862
if id == tool_use_id - 863
)) - 864
) - 865
} - 866
crate::types::EvidenceRef::CheckpointDiff { .. } - 867
| crate::types::EvidenceRef::FlowNode { .. } - 868
| crate::types::EvidenceRef::ChildSession { .. } - 869
| crate::types::EvidenceRef::ExternalOperation { .. } => false, - 870
}) - 871
} - 872
- 873
/// Reconcile managed items left running by a process restart. Only child - 874
/// sessions explicitly reported as live may remain running; every other - 875
/// running item is conservatively interrupted and made retry-eligible. - 876
/// Possible side effects are never replayed automatically. - 877
pub fn reconcile_running_work( - 878
&mut self, - 879
live_child_sessions: &std::collections::HashSet<String>, - 880
) -> Result<usize, SessionError> { - 881
self.reconcile_running_work_with_child_ledgers(live_child_sessions, None) - 882
} - 883
- 884
pub fn reconcile_running_work_with_child_ledgers( - 885
&mut self, - 886
live_child_sessions: &std::collections::HashSet<String>, - 887
sessions_home: Option<&Path>, - 888
) -> Result<usize, SessionError> { - 889
let Some(projection) = self - 890
.work_projection() - 891
.map_err(|error| SessionError::Corrupt { - 892
line: 0, - 893
message: error.to_string(), - 894
})? - 895
else { - 896
return Ok(0); - 897
}; - 898
let stale: Vec<(String, u32, Option<String>)> = projection - 899
.items - 900
.values() - 901
.filter(|item| { - 902
item.status == crate::types::WorkItemStatus::Running - 903
&& !item - 904
.child_session_id - 905
.as_ref() - 906
.is_some_and(|id| live_child_sessions.contains(id)) - 907
}) - 908
.map(|item| { - 909
( - 910
item.item_id.clone(), - 911
item.attempt, - 912
item.child_session_id.clone(), - 913
) - 914
}) - 915
.collect(); - 916
let mut reconciled = 0; - 917
for (item_id, attempt, child_id) in stale { - 918
let Some(current) = self - 919
.work_projection() - 920
.map_err(|error| SessionError::Corrupt { - 921
line: 0, - 922
message: error.to_string(), - 923
})? - 924
else { - 925
break; - 926
}; - 927
let child_status = child_id.as_deref().and_then(|id| { - 928
let home = sessions_home?; - 929
let cwd = self.header().map(|h| h.contract_cwd())?; - 930
let path = SessionPath::new_session_file(home, &cwd, id); - 931
SessionLog::open(path) - 932
.ok() - 933
.and_then(|child| child.child_run_status()) - 934
}); - 935
if matches!(child_status, Some(crate::types::ChildRunStatus::Completed)) { - 936
self.append_work(WorkEvent { - 937
contract_id: current.contract.contract_id.clone(), - 938
revision: current.contract.revision, - 939
kind: crate::types::WorkEventKind::EvidenceAttached { - 940
item_id: item_id.clone(), - 941
evidence: crate::types::EvidenceRef::ChildSession { - 942
session_id: child_id.clone().unwrap_or_default(), - 943
}, - 944
}, - 945
})?; - 946
self.append_work(WorkEvent { - 947
contract_id: current.contract.contract_id.clone(), - 948
revision: current.contract.revision, - 949
kind: crate::types::WorkEventKind::ItemStatusChanged { - 950
item_id, - 951
from: crate::types::WorkItemStatus::Running, - 952
to: crate::types::WorkItemStatus::ReadyForVerification, - 953
attempt, - 954
reason: "recovered completed child; verify its durable evidence".into(), - 955
}, - 956
})?; - 957
reconciled += 1; - 958
continue; - 959
} - 960
self.append_work(WorkEvent { - 961
contract_id: current.contract.contract_id.clone(), - 962
revision: current.contract.revision, - 963
kind: crate::types::WorkEventKind::ItemStatusChanged { - 964
item_id, - 965
from: crate::types::WorkItemStatus::Running, - 966
to: crate::types::WorkItemStatus::Interrupted, - 967
attempt, - 968
reason: match child_status { - 969
Some(crate::types::ChildRunStatus::Failed) => { - 970
"child failed before restart; review before retry" - 971
} - 972
Some(crate::types::ChildRunStatus::Aborted) => { - 973
"child was aborted before restart; review before retry" - 974
} - 975
Some(crate::types::ChildRunStatus::MaxTurns) => { - 976
"child hit its turn limit; review before retry" - 977
} - 978
_ => "recovered after process restart; review before retry", - 979
} - 980
.into(), - 981
}, - 982
})?; - 983
reconciled += 1; - 984
} - 985
Ok(reconciled) - 986
} - 987
- 988
pub fn activities( - 989
&self, - 990
) -> Vec<( - 991
String, - 992
chrono::DateTime<chrono::Utc>, - 993
crate::types::ActivityRecord, - 994
)> { - 995
self.chain_to_root() - 996
.into_iter() - 997
.filter_map(|entry| match &entry.payload { - 998
EntryPayload::Activity(activity) => { - 999
Some((entry.id.clone(), entry.ts, activity.clone())) - 1000
} - 1001
_ => None, - 1002
}) - 1003
.collect() - 1004
} - 1005
- 1006
/// Successful tool-result receipts with the ledger timestamp at which - 1007
/// the result was recorded. This is a replay-safe source for evidence - 1008
/// freshness; callers choose the domain-specific validity window. - 1009
pub fn successful_tool_receipts(&self) -> Vec<(String, chrono::DateTime<chrono::Utc>)> { - 1010
let mut known_calls = std::collections::HashSet::new(); - 1011
let mut receipts = Vec::new(); - 1012
for entry in self.chain_to_root() { - 1013
if let EntryPayload::Message(record) = &entry.payload { - 1014
for block in &record.message.content { - 1015
match block { - 1016
vak_llm::ContentBlock::ToolUse { id, .. } => { - 1017
known_calls.insert(id.clone()); - 1018
} - 1019
vak_llm::ContentBlock::ToolResult { - 1020
tool_use_id, - 1021
is_error: false, - 1022
.. - 1023
} if known_calls.contains(tool_use_id) => { - 1024
receipts.push((tool_use_id.clone(), entry.ts)); - 1025
} - 1026
_ => {} - 1027
} - 1028
} - 1029
} - 1030
} - 1031
receipts - 1032
} - 1033
- 1034
/// Successful receipts on the active, latest user turn only. Earlier - 1035
/// turns are deliberately excluded so a stale unrelated command cannot - 1036
/// establish evidence for the current result. - 1037
pub fn successful_tool_receipts_for_latest_turn( - 1038
&self, - 1039
) -> Vec<(String, chrono::DateTime<chrono::Utc>)> { - 1040
let mut known_calls = std::collections::HashSet::new(); - 1041
let mut receipts = Vec::new(); - 1042
for entry in self.chain_to_root() { - 1043
if let EntryPayload::Message(record) = &entry.payload { - 1044
if record.message.role == vak_llm::Role::User - 1045
&& record.control_kind().is_none() - 1046
&& record - 1047
.message - 1048
.content - 1049
.iter() - 1050
.any(|block| matches!(block, vak_llm::ContentBlock::Text { .. })) - 1051
{ - 1052
known_calls.clear(); - 1053
receipts.clear(); - 1054
} - 1055
for block in &record.message.content { - 1056
match block { - 1057
vak_llm::ContentBlock::ToolUse { id, .. } => { - 1058
known_calls.insert(id.clone()); - 1059
} - 1060
vak_llm::ContentBlock::ToolResult { - 1061
tool_use_id, - 1062
is_error: false, - 1063
.. - 1064
} if known_calls.contains(tool_use_id) => { - 1065
receipts.push((tool_use_id.clone(), entry.ts)); - 1066
} - 1067
_ => {} - 1068
} - 1069
} - 1070
} - 1071
} - 1072
receipts - 1073
} - 1074
- 1075
/// Distinct bash commands that ran GREEN on the active chain, in - 1076
/// first-run order (docs/design/10-flows.md adoption substrate). A command is - 1077
/// settled when its tool_result is not an error. - 1078
pub fn settled_bash_commands(&self) -> Vec<String> { - 1079
use std::collections::HashMap; - 1080
// id -> is_error for tool results - 1081
let mut results: HashMap<String, bool> = HashMap::new(); - 1082
for e in self.chain_to_root() { - 1083
if let EntryPayload::Message(r) = &e.payload { - 1084
for b in &r.message.content { - 1085
if let vak_llm::ContentBlock::ToolResult { - 1086
tool_use_id: id, - 1087
is_error, - 1088
.. - 1089
} = b - 1090
{ - 1091
results.insert(id.clone(), *is_error); - 1092
} - 1093
} - 1094
} - 1095
} - 1096
let mut out: Vec<String> = Vec::new(); - 1097
for e in self.chain_to_root() { - 1098
if let EntryPayload::Message(r) = &e.payload { - 1099
for b in &r.message.content { - 1100
if let vak_llm::ContentBlock::ToolUse { name, input, id } = b - 1101
&& name == "bash" - 1102
&& !results.get(id).copied().unwrap_or(true) - 1103
&& let Some(cmd) = input.get("command").and_then(|v| v.as_str()) - 1104
&& !out.iter().any(|o| o == cmd) - 1105
{ - 1106
out.push(cmd.to_string()); - 1107
} - 1108
} - 1109
} - 1110
} - 1111
out - 1112
} - 1113
- 1114
/// Receipt entries along the active path, root→leaf — forensic view - 1115
/// for surfaces that render dispatch history. - 1116
pub fn receipts(&self) -> Vec<&vak_llm::WorkReceipt> { - 1117
self.chain_to_root() - 1118
.into_iter() - 1119
.filter_map(|e| match &e.payload { - 1120
EntryPayload::Receipt(r) => Some(r), - 1121
_ => None, - 1122
}) - 1123
.collect() - 1124
} - 1125
- 1126
pub fn branch_at(&mut self, entry_id: &str) -> Result<(), SessionError> { - 1127
if !self.by_id.contains_key(entry_id) { - 1128
return Err(SessionError::Corrupt { - 1129
line: 0, - 1130
message: format!("cannot branch at unknown entry {entry_id}"), - 1131
}); - 1132
} - 1133
self.tail_id = Some(entry_id.to_string()); - 1134
Ok(()) - 1135
} - 1136
- 1137
/// Appends one compaction packet over the inclusive turn range - 1138
/// `first_turn_id..=last_turn_id` (docs/design/68-context-engine.md §4). - 1139
/// Both ids must be directive entries on the chain. The packet is a - 1140
/// cache keyed by that range: it never moves a boundary and never hides - 1141
/// the turns it covers from a plan that wants them at `Full` or `Card`. - 1142
pub fn append_packet( - 1143
&mut self, - 1144
first_turn_id: &str, - 1145
last_turn_id: &str, - 1146
model: &str, - 1147
summary: String, - 1148
tokens_before: u64, - 1149
) -> Result<Entry, SessionError> { - 1150
for id in [first_turn_id, last_turn_id] { - 1151
if !self.by_id.contains_key(id) { - 1152
return Err(SessionError::Corrupt { - 1153
line: 0, - 1154
message: format!("unknown turn id {id} in packet range"), - 1155
}); - 1156
} - 1157
} - 1158
let parent = self.tail_id.clone(); - 1159
self.append(Entry::new( - 1160
parent, - 1161
EntryPayload::Compaction(crate::types::CompactionEntry { - 1162
summary, - 1163
first_turn_id: first_turn_id.to_string(), - 1164
last_turn_id: last_turn_id.to_string(), - 1165
model: model.to_string(), - 1166
tokens_before, - 1167
reset_all: false, - 1168
}), - 1169
)) - 1170
} - 1171
- 1172
/// Reset-with-handoff (docs/design/42-managed-work-contracts.md): the projection becomes ONLY - 1173
/// this summary. Append-only; the full history stays on disk. This is - 1174
/// the one compaction entry that is a real boundary — the rescue for a - 1175
/// profile with no usable horizon, where nothing is plannable. - 1176
pub fn append_handoff_reset( - 1177
&mut self, - 1178
summary: String, - 1179
tokens_before: u64, - 1180
) -> Result<Entry, SessionError> { - 1181
let parent = self.tail_id.clone(); - 1182
self.append(Entry::new( - 1183
parent, - 1184
EntryPayload::Compaction(crate::types::CompactionEntry { - 1185
summary, - 1186
first_turn_id: String::new(), - 1187
last_turn_id: String::new(), - 1188
model: String::new(), - 1189
tokens_before, - 1190
reset_all: true, - 1191
}), - 1192
)) - 1193
} - 1194
- 1195
pub fn header(&self) -> Option<&SessionHeader> { - 1196
self.entries.iter().find_map(|e| match &e.payload { - 1197
EntryPayload::Header(h) => Some(h), - 1198
_ => None, - 1199
}) - 1200
} - 1201
- 1202
pub fn len(&self) -> usize { - 1203
self.entries.len() - 1204
} - 1205
- 1206
pub fn is_empty(&self) -> bool { - 1207
self.entries.is_empty() - 1208
} - 1209
- 1210
pub fn path(&self) -> &Path { - 1211
&self.path - 1212
} - 1213
- 1214
pub fn tail_id(&self) -> Option<&String> { - 1215
self.tail_id.as_ref() - 1216
} - 1217
- 1218
pub fn chain_to_root(&self) -> Vec<&Entry> { - 1219
let mut chain = Vec::new(); - 1220
let mut cursor = self.tail_id.clone(); - 1221
while let Some(id) = cursor { - 1222
let Some(&idx) = self.by_id.get(&id) else { - 1223
break; - 1224
}; - 1225
let entry = &self.entries[idx]; - 1226
chain.push(entry); - 1227
cursor = entry.parent_id.clone(); - 1228
} - 1229
chain.reverse(); - 1230
chain - 1231
} - 1232
- 1233
/// Message entries along the active path, root→leaf, with their entry - 1234
/// ids — raw ledger view (compaction entries NOT applied). - 1235
pub fn message_chain(&self) -> Vec<(String, Message)> { - 1236
self.chain_to_root() - 1237
.into_iter() - 1238
.filter_map(|e| match &e.payload { - 1239
EntryPayload::Message(r) => Some((e.id.clone(), r.message.clone())), - 1240
_ => None, - 1241
}) - 1242
.collect() - 1243
} - 1244
- 1245
fn derive_keyed(&self) -> Vec<(String, Message)> { - 1246
self.derive_keyed_tagged() - 1247
.into_iter() - 1248
.map(|(id, m, _, _)| (id, m)) - 1249
.collect() - 1250
} - 1251
- 1252
/// The plan-free projection: every closed turn at `Full` fidelity. Used - 1253
/// by `derive_messages` (goal audits, search, human-facing views), which - 1254
/// has no `CapacityProfile` to plan against. - 1255
fn derive_keyed_tagged(&self) -> Vec<(String, Message, bool, bool)> { - 1256
self.derive_with_plan_tagged(None) - 1257
} - 1258
- 1259
/// Like `derive_keyed_tagged`, but a `WorkingSetPlan` (from - 1260
/// `vak_context::planner::plan`) selects each closed turn's fidelity - 1261
/// instead of defaulting every one to `Full` (docs/design/68-context- - 1262
/// engine.md §4/§10) — one implementation, the plan just picks what each - 1263
/// turn contributes: - 1264
/// - 1265
/// - `Full` turns project their `full_record`: real `tool_use` blocks, - 1266
/// results as schema-driven digests carrying their evidence id, never - 1267
/// a character-count trim. - 1268
/// - `Card` turns contribute one line each to a single `<turns>` block - 1269
/// (`TurnCard::line`), inserted once, right after any compaction - 1270
/// summary. - 1271
/// - `Packet` turns, and any closed turn the plan omits entirely, - 1272
/// contribute nothing here: they are represented only by an existing - 1273
/// `Compaction` entry, which the caller (`Agent`'s incremental - 1274
/// compaction, §4) guarantees already covers them before this is - 1275
/// called with that plan. - 1276
/// - The still-open turn (if any) always projects verbatim, regardless - 1277
/// of `plan` — it is never planned. - 1278
/// - 1279
/// The intent note, work contract, and conversation thread are rendered - 1280
/// into the request tail instead (§6/§10), read separately via - 1281
/// [`SessionLog::tail_sections`]; this function never contributes them. - 1282
/// The reset boundary: the chain position before which everything is - 1283
/// invisible to the model, and the handoff summary that stands in for - 1284
/// it. Only a `reset_all` compaction entry (reset-with-handoff, - 1285
/// docs/design/42) moves this; packets never do. `(0, ..., None)` when - 1286
/// no reset has happened. - 1287
fn reset_boundary(&self) -> (usize, HashMap<String, usize>, Option<String>) { - 1288
let chain = self.chain_to_root(); - 1289
let position: HashMap<String, usize> = chain - 1290
.iter() - 1291
.enumerate() - 1292
.map(|(i, entry)| (entry.id.clone(), i)) - 1293
.collect(); - 1294
let last_reset = - 1295
chain - 1296
.iter() - 1297
.enumerate() - 1298
.rev() - 1299
.find_map(|(pos, entry)| match &entry.payload { - 1300
EntryPayload::Compaction(c) if c.reset_all => Some((pos, c.summary.clone())), - 1301
_ => None, - 1302
}); - 1303
match last_reset { - 1304
Some((pos, summary)) => (pos, position, Some(summary)), - 1305
None => (0, position, None), - 1306
} - 1307
} - 1308
- 1309
/// Every stored packet (non-reset compaction entry) in ledger order. - 1310
fn packets(&self) -> Vec<Packet> { - 1311
TurnIndex::from_log(self).packets - 1312
} - 1313
- 1314
/// The stored packet whose range is exactly `first_turn_id..=last_turn_id`, - 1315
/// newest such entry first (docs/design/68 §4: a packet is reused only - 1316
/// for the exact range the plan asks for). `None` when no packet - 1317
/// covers that range, whatever other packets exist. - 1318
pub fn packet_for(&self, first_turn_id: &str, last_turn_id: &str) -> Option<Packet> { - 1319
self.packets().into_iter().rev().find(|packet| { - 1320
packet.first_turn_id == first_turn_id && packet.last_turn_id == last_turn_id - 1321
}) - 1322
} - 1323
- 1324
/// Whether the packet range a `WorkingSetPlan` asked for - 1325
/// (`packet_range = (first, last)`) still needs a summariser call: true - 1326
/// when no stored packet covers exactly that range. - 1327
pub fn packet_needs_compaction(&self, first_turn_id: &str, last_turn_id: &str) -> bool { - 1328
self.packet_for(first_turn_id, last_turn_id).is_none() - 1329
} - 1330
- 1331
/// The summariser input for a packet over `first_turn_id..=last_turn_id`: - 1332
/// the longest stored packet that starts at the same turn and ends at or - 1333
/// before `last_turn_id` (so work already folded in is reused, never - 1334
/// re-read from raw history), followed by one `TurnCard` line per turn - 1335
/// after it up to `last_turn_id` — cards, never raw history (§4). Also - 1336
/// returns a char-count estimate for the caller's `tokens_before`. - 1337
pub fn packet_transcript(&self, first_turn_id: &str, last_turn_id: &str) -> (String, u64) { - 1338
let (_, position, _) = self.reset_boundary(); - 1339
let mut out = String::new(); - 1340
let (Some(&first_pos), Some(&last_pos)) = - 1341
(position.get(first_turn_id), position.get(last_turn_id)) - 1342
else { - 1343
return (out, 0); - 1344
}; - 1345
// Seed: the stored packet with this `first` and the greatest `last` - 1346
// not past our target. - 1347
let seed = self - 1348
.packets() - 1349
.into_iter() - 1350
.filter(|packet| packet.first_turn_id == first_turn_id) - 1351
.filter_map(|packet| { - 1352
let end = position.get(packet.last_turn_id.as_str()).copied()?; - 1353
(end <= last_pos).then_some((end, packet)) - 1354
}) - 1355
.max_by_key(|(end, _)| *end); - 1356
let mut from_pos = first_pos; - 1357
if let Some((end, packet)) = seed { - 1358
out.push_str(&packet.summary); - 1359
out.push_str("\n\n"); - 1360
from_pos = end + 1; - 1361
} - 1362
let index = TurnIndex::from_log(self); - 1363
for (turn_number, turn) in index.turns.iter().enumerate() { - 1364
let turn_pos = position.get(turn.id.as_str()).copied().unwrap_or(0); - 1365
if turn_pos < from_pos || turn_pos > last_pos { - 1366
continue; - 1367
} - 1368
match &turn.card { - 1369
Some(card) => { - 1370
out.push_str(&card.line(turn_number + 1)); - 1371
out.push('\n'); - 1372
} - 1373
None => { - 1374
for message in turn.full_record() { - 1375
out.push_str(&message.text_content()); - 1376
out.push('\n'); - 1377
} - 1378
} - 1379
} - 1380
} - 1381
let chars = out.chars().count() as u64; - 1382
(out, chars) - 1383
} - 1384
- 1385
/// The model-visible projection for one request. With a plan, each - 1386
/// closed turn contributes exactly what the plan chose for it; without - 1387
/// one, every closed turn rides at `Full` (the plan-free view used by - 1388
/// goal audits and search, which have no `CapacityProfile`). Only the - 1389
/// reset boundary (reset-with-handoff) hides anything; a stored packet - 1390
/// is rendered only when the plan's `packet_range` matches it exactly, - 1391
/// and a plan that packets a range no stored packet covers yet renders - 1392
/// those turns as cards — the cheap, safe form — rather than losing - 1393
/// them (no-cut invariant). The caller normally writes the packet - 1394
/// before projecting (`Agent`'s incremental compaction, §4), so that - 1395
/// fallback is a transient. - 1396
fn derive_with_plan_tagged( - 1397
&self, - 1398
plan: Option<&WorkingSetPlan>, - 1399
) -> Vec<(String, Message, bool, bool)> { - 1400
let (boundary_pos, position_owned, reset_summary) = self.reset_boundary(); - 1401
let position: HashMap<&str, usize> = position_owned - 1402
.iter() - 1403
.map(|(id, pos)| (id.as_str(), *pos)) - 1404
.collect(); - 1405
- 1406
let mut index = TurnIndex::from_log(self); - 1407
// Rendering needs a card line for every closed turn the plan put at - 1408
// `Card`; the token counts on a provisional card are irrelevant - 1409
// here (the plan was costed by the caller), so they are left at 0. - 1410
index.ensure_cards(&|_| 0); - 1411
let mut out: Vec<(String, Message, bool, bool)> = Vec::new(); - 1412
if let Some(summary) = &reset_summary { - 1413
let summary_msg = - 1414
Message::user_text(format!("<context_summary>\n{summary}\n</context_summary>")); - 1415
out.push((String::new(), summary_msg, true, false)); - 1416
} - 1417
- 1418
let fidelity_of: HashMap<&str, Fidelity> = plan - 1419
.map(|p| p.per_turn.iter().map(|(id, f)| (id.as_str(), *f)).collect()) - 1420
.unwrap_or_default(); - 1421
// `plan.per_turn` never carries `Fidelity::Packet` entries (§4): a - 1422
// packeted turn is identified by falling inside `packet_range` - 1423
// instead. Resolved once, by position, so membership is an O(1) - 1424
// range check per turn. - 1425
let packet_pos_range: Option<(usize, usize)> = plan.and_then(|p| { - 1426
let (first, last) = p.packet_range.as_ref()?; - 1427
let lo = position.get(first.as_str()).copied()?; - 1428
let hi = position.get(last.as_str()).copied()?; - 1429
Some((lo.min(hi), lo.max(hi))) - 1430
}); - 1431
let packet = plan - 1432
.and_then(|p| p.packet_range.as_ref()) - 1433
.and_then(|(first, last)| self.packet_for(first, last)); - 1434
if let Some(packet) = &packet { - 1435
let summary_msg = Message::user_text(format!( - 1436
"<context_summary>\n{}\n</context_summary>", - 1437
packet.summary - 1438
)); - 1439
out.push((String::new(), summary_msg, true, false)); - 1440
} - 1441
- 1442
let mut card_lines: Vec<String> = Vec::new(); - 1443
for (turn_number, turn) in index.turns.iter().enumerate() { - 1444
let turn_pos = position.get(turn.id.as_str()).copied().unwrap_or(0); - 1445
if turn_pos < boundary_pos { - 1446
continue; - 1447
} - 1448
if !turn.closed { - 1449
for message in turn.current_verbatim() { - 1450
out.push((turn.id.clone(), message, false, false)); - 1451
} - 1452
continue; - 1453
} - 1454
let selected = match plan { - 1455
None => Fidelity::Full, - 1456
Some(_) => { - 1457
if let Some(f) = fidelity_of.get(turn.id.as_str()).copied() { - 1458
f - 1459
} else if packet_pos_range - 1460
.is_some_and(|(lo, hi)| turn_pos >= lo && turn_pos <= hi) - 1461
{ - 1462
if packet.is_some() { - 1463
Fidelity::Packet - 1464
} else { - 1465
// The plan asked for a packet nobody has - 1466
// written yet: cards, never nothing. - 1467
Fidelity::Card - 1468
} - 1469
} else { - 1470
// The plan never classified it (a stale plan - 1471
// against a longer chain) — render the cheap, safe - 1472
// form rather than silently losing it (no-cut - 1473
// invariant); every closed turn has a card here. - 1474
Fidelity::Card - 1475
} - 1476
} - 1477
}; - 1478
match selected { - 1479
Fidelity::Full => { - 1480
for message in turn.full_record() { - 1481
out.push((turn.id.clone(), message, false, false)); - 1482
} - 1483
} - 1484
Fidelity::Card => { - 1485
if let Some(card) = &turn.card { - 1486
card_lines.push(card.line(turn_number + 1)); - 1487
} - 1488
} - 1489
Fidelity::Packet => { - 1490
// Represented by the packet summary pushed above. - 1491
} - 1492
} - 1493
} - 1494
if !card_lines.is_empty() { - 1495
let block = format!("<turns>\n{}\n</turns>", card_lines.join("\n")); - 1496
let insert_at = usize::from(reset_summary.is_some()) + usize::from(packet.is_some()); - 1497
out.insert( - 1498
insert_at, - 1499
(String::new(), Message::user_text(block), true, false), - 1500
); - 1501
} - 1502
out - 1503
} - 1504
- 1505
/// `derive_keyed_tagged` with a `WorkingSetPlan` applied — the projection - 1506
/// the request assembler sends once a `CapacityProfile` is available - 1507
/// (docs/design/68-context-engine.md §4/§10). - 1508
pub fn derive_with_plan(&self, plan: &WorkingSetPlan) -> Vec<Message> { - 1509
self.derive_with_plan_tagged(Some(plan)) - 1510
.into_iter() - 1511
.map(|(_, m, _, _)| m) - 1512
.collect() - 1513
} - 1514
- 1515
/// `derive_with_plan`'s projection alongside the index within it where - 1516
/// the open turn's own directive sits — `messages.len()` (an - 1517
/// out-of-bounds sentinel, safe as a no-op) when the chain has no open - 1518
/// turn. What a request assembler needs to attach the per-turn tail to - 1519
/// the turn's directive specifically (docs/design/68-context-engine.md - 1520
/// §6/§7): the directive is always the first message of the open - 1521
/// turn's own `current_verbatim` block, which `derive_with_plan` - 1522
/// places last, but a step within the turn appends more messages - 1523
/// (tool results, control nudges) after it — the tail must ride on the - 1524
/// directive every step, not on whatever the last message happens to - 1525
/// be, or an earlier-sent message would silently change shape between - 1526
/// requests. - 1527
pub fn derive_with_plan_and_directive(&self, plan: &WorkingSetPlan) -> (Vec<Message>, usize) { - 1528
let messages = self.derive_with_plan(plan); - 1529
let index = TurnIndex::from_log(self); - 1530
let directive_at = index - 1531
.turns - 1532
.last() - 1533
.filter(|turn| !turn.closed) - 1534
.map(|turn| messages.len().saturating_sub(turn.current_verbatim().len())) - 1535
.unwrap_or(messages.len()); - 1536
(messages, directive_at) - 1537
} - 1538
- 1539
/// Appends the packet a `WorkingSetPlan` asked for (its `packet_range`, - 1540
/// `first_turn_id..=last_turn_id`), produced by INCREMENTAL compaction - 1541
/// (§4): summarizing the turns' CARDS (never raw history). `model` is - 1542
/// the model whose plan asked for it. The packet is keyed by its range - 1543
/// and is reused by any later plan — under any model — that asks for - 1544
/// exactly that range; it hides nothing from a plan that does not. - 1545
pub fn append_incremental_compaction( - 1546
&mut self, - 1547
first_turn_id: &str, - 1548
last_turn_id: &str, - 1549
model: &str, - 1550
summary: String, - 1551
tokens_before: u64, - 1552
) -> Result<(), SessionError> { - 1553
self.append_packet(first_turn_id, last_turn_id, model, summary, tokens_before)?; - 1554
Ok(()) - 1555
} - 1556
- 1557
/// Every raw message entry in chain order, tagged with its ledger - 1558
/// identity and class — a human-facing audit view, independent of the - 1559
/// turn-based model-visible projection (`derive_messages`). Unlike that - 1560
/// projection, this one is never affected by compaction: a person - 1561
/// reading their own history sees everything they said, not what a - 1562
/// context budget kept. Compaction entries still surface as a - 1563
/// `context: true` pseudo-message so a consumer can show "context - 1564
/// summarized here" inline. - 1565
pub fn derive_transcript(&self) -> Vec<TranscriptMessage> { - 1566
self.chain_to_root() - 1567
.into_iter() - 1568
.filter_map(|entry| match &entry.payload { - 1569
EntryPayload::Message(record) => Some(TranscriptMessage { - 1570
entry_id: entry.id.clone(), - 1571
message: record.message.clone(), - 1572
control: record.control_kind(), - 1573
context: false, - 1574
author_id: record.meta.as_ref().and_then(|meta| meta.author_id.clone()), - 1575
author_name: record - 1576
.meta - 1577
.as_ref() - 1578
.and_then(|meta| meta.author_name.clone()), - 1579
attachments: record - 1580
.meta - 1581
.as_ref() - 1582
.map(|meta| meta.attachments.clone()) - 1583
.unwrap_or_default(), - 1584
}), - 1585
EntryPayload::Compaction(c) => Some(TranscriptMessage { - 1586
entry_id: entry.id.clone(), - 1587
message: Message::user_text(format!( - 1588
"<context_summary>\n{}\n</context_summary>", - 1589
c.summary - 1590
)), - 1591
control: None, - 1592
context: true, - 1593
author_id: None, - 1594
author_name: None, - 1595
attachments: Vec::new(), - 1596
}), - 1597
_ => None, - 1598
}) - 1599
.collect() - 1600
} - 1601
- 1602
/// The conversation as a person would read it: the model-visible messages - 1603
/// minus the runtime-authored nudges. For exports and other human-facing - 1604
/// views; the model itself is fed `derive_messages`, which is unchanged. - 1605
pub fn derive_conversation(&self) -> Vec<Message> { - 1606
self.derive_transcript() - 1607
.into_iter() - 1608
.filter(|item| item.control.is_none()) - 1609
.map(|item| item.message) - 1610
.collect() - 1611
} - 1612
- 1613
pub fn derive_messages(&self) -> Vec<Message> { - 1614
self.derive_keyed().into_iter().map(|(_, m)| m).collect() - 1615
} - 1616
- 1617
/// The three sections `derive_keyed_tagged` used to splice into the - 1618
/// model-visible projection, computed separately for the request - 1619
/// assembler's tail (docs/design/68-context-engine.md §6/§10). Each - 1620
/// field is `None` when there is nothing to say — the caller renders no - 1621
/// tag for an absent section, never an empty one. - 1622
/// The current directive's resolved reading, distilled to a - 1623
/// `ReadingKey` (docs/design/68-context-engine.md §4's `PlanInput`) — - 1624
/// the newest `Intent` entry on the chain, whether or not its turn has - 1625
/// closed yet. `None` before the first intent reading of the session. - 1626
/// The still-open turn's own verbatim messages, respecting the reset - 1627
/// boundary (docs/design/68-context-engine.md §4): after a - 1628
/// reset-with-handoff, everything before the reset entry is invisible - 1629
/// to the model even though the open turn technically started before - 1630
/// it, so planning must size its reserve against what will actually be - 1631
/// sent — never against text the reset already discarded. Empty when - 1632
/// there is no open turn, or the open turn itself predates the - 1633
/// boundary. - 1634
pub fn open_turn_verbatim(&self) -> Vec<Message> { - 1635
let (boundary_pos, position, _) = self.reset_boundary(); - 1636
let index = TurnIndex::from_log(self); - 1637
let Some(turn) = index.turns.last().filter(|t| !t.closed) else { - 1638
return Vec::new(); - 1639
}; - 1640
let turn_pos = position.get(turn.id.as_str()).copied().unwrap_or(0); - 1641
if turn_pos < boundary_pos { - 1642
return Vec::new(); - 1643
} - 1644
turn.current_verbatim() - 1645
} - 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
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.