fix(workstream): harden Kimi managed continuity

This commit is contained in:
AkitaOnRails
2026-07-23 11:57:26 -03:00
parent acd424f3fb
commit 260af15c8d
13 changed files with 652 additions and 104 deletions
+18 -4
View File
@@ -44,14 +44,28 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
model. `UserPromptSubmit` stdout is injected as a user message before the
turn. The native `ai-memory hook` path now also accepts the
`user-prompt-submit` event token alongside `user-prompt` for kimi handoff
delivery. Re-run `ai-memory install-hooks --agent kimi-code --apply` after
upgrading.
delivery. Existing native hook commands pick up the correction when the
binary is upgraded; script-fallback installations must re-run
`ai-memory install-hooks --agent kimi-code --apply` to refresh their staged
scripts.
- The `[briefing]` compiled project brief is now gated to once per session
for Kimi Code: it rides the first user prompt (kimi discards SessionStart
hook stdout, so session-start delivery is impossible) and later prompts
keep fetching the handoff without the briefing params, so the server no
longer recomposes the brief on every prompt. Re-briefing after `/clear`
is not supported in v1.
longer recomposes the brief on every prompt. This also works for managed
handoffs, creates local markers only for opted-in repositories, and retains
at most 512 markers. Re-briefing after `/clear` is not supported in v1.
- Kimi Code managed transcript cursors now validate the imported byte prefix
before resuming. When Kimi rewrites `wire.jsonl` in place, ai-memory safely
replays the journal with stable event IDs instead of seeking into a changed
record or skipping new events.
- Kimi Code managed transcript extraction now imports native
`context.append_loop_event` assistant text, tool calls, and tool results.
Current Kimi journals store model output in those records rather than
re-appending it as a completed assistant message.
- `ai-memory run kimi server ...` now passes through Kimi Code's deprecated
but still functional `server` utility command instead of treating it as an
interactive session launch.
## [1.17.3] - 2026-07-22
+104 -22
View File
@@ -27,7 +27,8 @@ use sha2::{Digest as _, Sha256};
use super::hook_capture::{
build_client, canonical_context, capture_policy, extract_cwd, get_handoff, marker_query_suffix,
marker_query_suffix_without_briefing, resolve_cwd_with_fallbacks, url_encode,
marker_query_suffix_without_briefing, marker_requests_briefing, resolve_cwd_with_fallbacks,
url_encode,
};
use super::hook_drain_process;
use super::hook_spool;
@@ -54,6 +55,7 @@ const MANAGED_RUN_ENV: &str = "AI_MEMORY_RUN_ID";
/// Backlog size at which `post-tool-use` does a mid-session catch-up drain, so a
/// light session pays only a `read_dir`. Override via the env var above.
const DEFAULT_INCREMENTAL_THRESHOLD: usize = 32;
const MAX_BRIEFED_MARKERS: usize = 512;
/// Total budget AND per-event timeout for the mid-session catch-up drain — kept
/// well under a second so a `post-tool-use` hook never stalls a tool call (one
/// in-flight POST against a slow server is bounded by this too).
@@ -159,9 +161,9 @@ fn clear_session_id(data_dir: &Path, agent: AgentKind) {
/// hook stdout, so the brief rides the FIRST user prompt of the session
/// (parity with Claude's once-per-SessionStart brief); the marker keeps
/// later prompts from re-requesting it. Keyed by the payload's canonical
/// session id (kimi always sends `sessionId`); payloads without one fall
/// back to a stable hash of agent+cwd so a session-less agent still briefs
/// once per checkout. The key is sanitized to a safe file name.
/// session id when Kimi supplies one; payloads without one fall back to a
/// stable hash of agent+cwd so a session-less agent still briefs once per
/// checkout. The key is sanitized to a safe file name.
fn briefed_marker_path(
data_dir: &Path,
agent: &str,
@@ -192,13 +194,36 @@ fn sanitize_briefed_key(raw: &str) -> String {
.collect()
}
/// Best-effort marker write: on failure the worst case is a re-brief on the
/// next prompt, which is acceptable.
/// Best-effort marker write with bounded retention: on failure the worst case
/// is a re-brief on the next prompt, which is acceptable.
fn mark_briefed(path: &Path) {
if let Some(parent) = path.parent() {
let _ = fs::create_dir_all(parent);
let Some(parent) = path.parent() else {
return;
};
if fs::create_dir_all(parent).is_err() || fs::write(path, b"").is_err() {
return;
}
let Ok(entries) = fs::read_dir(parent) else {
return;
};
let mut stale_candidates = entries
.filter_map(Result::ok)
.filter_map(|entry| {
let metadata = entry.metadata().ok()?;
metadata.is_file().then_some((
metadata.modified().ok(),
entry.file_name(),
entry.path(),
))
})
.filter(|(_, _, candidate)| candidate != path)
.collect::<Vec<_>>();
stale_candidates.sort_by(|a, b| b.0.cmp(&a.0).then_with(|| b.1.cmp(&a.1)));
let keep_others = MAX_BRIEFED_MARKERS.saturating_sub(1);
for (_, _, stale) in stale_candidates.into_iter().skip(keep_others) {
let _ = fs::remove_file(stale);
}
let _ = fs::write(path, b"");
}
fn fresh_session_id(data_dir: &Path, agent: AgentKind) -> String {
@@ -566,8 +591,9 @@ where
// user-prompt: agents whose SessionStart stdout is discarded (Kimi Code)
// receive the handoff here instead — kimi injects UserPromptSubmit stdout
// into the turn verbatim as a `hook_result` user message. The payload
// carries the native `sessionId`, so the destructive GET can also link
// the managed run to the native session, same as session-start does.
// carries the native session id when available, so the destructive GET
// can also link the managed run to the native session, same as
// session-start does.
// The installed kimi hook passes the script stem (`user-prompt-submit`)
// while the legacy shell path posts `user-prompt`; HookEvent::parse
// canonicalizes both (and the snake/native spellings) to UserPrompt.
@@ -589,13 +615,18 @@ where
// `&briefing`/`&briefing_budget`, so the server does not recompose
// the brief per prompt. The marker survives `/clear`, so re-briefing
// after a context clear is not supported in v1.
let briefed_path = briefed_marker_path(
&dd,
&args.agent,
canonical_session_id.as_deref(),
policy_cwd.as_deref(),
);
let handoff_qs = if briefed_path.is_file() {
let briefed_path = policy_cwd
.as_deref()
.filter(|cwd| marker_requests_briefing(cwd))
.map(|_| {
briefed_marker_path(
&dd,
&args.agent,
canonical_session_id.as_deref(),
policy_cwd.as_deref(),
)
});
let handoff_qs = if briefed_path.as_ref().is_some_and(|path| path.is_file()) {
policy_cwd
.as_deref()
.map(|cwd| {
@@ -619,7 +650,9 @@ where
// the brief-flagged request on every prompt would not deliver
// anything anyway, and the one lost brief is recovered on the next
// session.
mark_briefed(&briefed_path);
if let Some(path) = briefed_path.as_deref() {
mark_briefed(path);
}
if let Some(handoff) = handoff {
writeln!(stdout, "{handoff}")?;
}
@@ -1552,6 +1585,7 @@ mod tests {
auth_token: None,
project_strategy: None,
check_capture: false,
capture_assistant: false,
}
}
@@ -1816,9 +1850,9 @@ mod tests {
std::fs::create_dir(&cwd).unwrap();
write_briefing_marker(&cwd);
let (base, mut requests) = serve_requests("200 OK", "AMWS-HANDOFF-DELTA").await;
// No sessionId in the payload (kimi always sends one; this is the
// defensive path): the briefed marker is keyed by a stable hash of
// agent+cwd, so a session-less payload still briefs only once.
// No session id in the payload: the briefed marker is keyed by a
// stable hash of agent+cwd, so a session-less payload still briefs
// only once.
let payload = serde_json::json!({"cwd": cwd, "prompt": "hi"}).to_string();
let mut stdout = Vec::new();
@@ -1855,4 +1889,52 @@ mod tests {
assert!(second.starts_with("GET /handoff?"), "{second}");
assert!(!second.contains("briefing"), "{second}");
}
#[tokio::test]
async fn kimi_user_prompt_without_briefing_creates_no_marker() {
let tmp = tempfile::tempdir().unwrap();
let data_dir = tmp.path().join("data");
let cwd = tmp.path().join("repo");
std::fs::create_dir(&cwd).unwrap();
let (base, mut requests) = serve_requests("404 Not Found", "").await;
let mut stdout = Vec::new();
run_with_payload(
Some(data_dir.clone()),
kimi_hook_args("user-prompt-submit", &base),
serde_json::json!({
"session_id": "session_abc",
"cwd": cwd,
"prompt": "hi"
})
.to_string(),
&mut stdout,
|_| Ok(()),
)
.await
.unwrap();
let request = first_request(&mut requests).await.unwrap();
assert!(request.starts_with("GET /handoff?"), "{request}");
assert!(!request.contains("briefing"), "{request}");
assert!(!data_dir.join("briefed").exists());
}
#[test]
fn briefed_markers_are_bounded_and_keep_current() {
let tmp = tempfile::tempdir().unwrap();
let marker_dir = tmp.path().join("briefed");
let count = MAX_BRIEFED_MARKERS + 20;
for index in 0..count {
mark_briefed(&marker_dir.join(format!("session-{index:04}")));
}
let retained = std::fs::read_dir(&marker_dir).unwrap().count();
assert_eq!(retained, MAX_BRIEFED_MARKERS);
assert!(
marker_dir
.join(format!("session-{:04}", count - 1))
.is_file()
);
}
}
@@ -224,6 +224,16 @@ pub fn marker_query_suffix_without_briefing(cwd: &str, default_strategy: Option<
marker_query_suffix_impl(cwd, default_strategy, false)
}
/// Whether the nearest marker explicitly enables the compiled project brief.
///
/// Kimi Code uses this before creating its local once-per-session marker so
/// repositories that did not opt in do not accumulate marker files.
pub fn marker_requests_briefing(cwd: &str) -> bool {
find_marker(cwd)
.and_then(|marker| parse_toml_flag(&marker, "inject_on_session_start"))
.is_some_and(|value| is_truthy(&value))
}
fn marker_query_suffix_impl(
cwd: &str,
default_strategy: Option<&str>,
@@ -319,6 +329,13 @@ fn parse_toml_flag(file: &Path, key: &str) -> Option<String> {
None
}
fn is_truthy(value: &str) -> bool {
matches!(
value.trim().to_ascii_lowercase().as_str(),
"1" | "true" | "yes" | "on"
)
}
fn repo_root_project(cwd: &str) -> Option<String> {
let root = ai_memory_consolidate::discover_main_repo_root(Path::new(cwd)).ok()?;
root.file_name()
@@ -928,6 +945,27 @@ drop_subagent_captures = "true"
assert!(!qs.contains("briefing"), "{qs}");
}
#[test]
fn marker_requests_briefing_only_for_truthy_opt_in() {
let tmp = tempfile::TempDir::new().unwrap();
let marker = tmp.path().join(".ai-memory.toml");
let cwd = tmp.path().to_str().unwrap();
assert!(!marker_requests_briefing(cwd));
std::fs::write(
&marker,
"[briefing]\ninject_on_session_start = false\nmax_chars = 6000\n",
)
.unwrap();
assert!(!marker_requests_briefing(cwd));
std::fs::write(
&marker,
"[briefing]\ninject_on_session_start = YeS\nmax_chars = 6000\n",
)
.unwrap();
assert!(marker_requests_briefing(cwd));
}
#[test]
fn marker_query_suffix_repo_root_non_git_keeps_project_implicit() {
let tmp = tempfile::TempDir::new().unwrap();
+178 -20
View File
@@ -1068,8 +1068,10 @@ async fn fetch_and_accept_handoff(
);
return Ok(None);
}
let brief_md =
resolve_requested_session_brief(state, &query, actor_user.as_deref()).await?;
if context.context_delivered {
return Ok(None);
return Ok(brief_md);
}
let rendered = render_managed_context(
&context.events,
@@ -1078,7 +1080,7 @@ async fn fetch_and_accept_handoff(
context.sync_after,
);
let _ = state.writer.accept_managed_run_context(run_id).await?;
return Ok(rendered);
return Ok(combine_handoff_and_brief(rendered, brief_md));
}
// `/handoff` has no session_id in the request — `per_session` mode
// therefore falls back to the single slot (graceful degradation),
@@ -1102,7 +1104,7 @@ async fn fetch_and_accept_handoff(
let handoff_md = {
let handoff = state
.reader
.latest_open_handoff(ws, proj, query.cwd)
.latest_open_handoff(ws, proj, query.cwd.clone())
.await?;
match handoff {
Some(h) => {
@@ -1115,27 +1117,72 @@ async fn fetch_and_accept_handoff(
// The brief is additive and non-destructive: unlike the handoff (a
// single-use slot consumed above), it is recomposed on every opted-in
// session start — exactly what a Claude Code `/clear` needs (#176).
let brief_md = if crate::payload::query_flag_truthy(query.briefing.as_deref()) {
let budget = query
.briefing_budget
.as_deref()
.and_then(|v| v.trim().parse::<usize>().ok())
.unwrap_or(BRIEF_BUDGET_DEFAULT)
.clamp(BRIEF_BUDGET_MIN, BRIEF_BUDGET_MAX);
let (core, recent) = state
.reader
.session_brief_pages(ws, proj, BRIEF_CORE_PAGES_LIMIT, BRIEF_RECENT_PAGES_LIMIT)
.await?;
render_session_brief(&core, &recent, budget)
} else {
None
let brief_md = render_requested_session_brief(state, &query, ws, proj).await?;
Ok(combine_handoff_and_brief(handoff_md, brief_md))
}
async fn resolve_requested_session_brief(
state: &HookState,
query: &HandoffQuery,
actor_user: Option<&str>,
) -> anyhow::Result<Option<String>> {
if !crate::payload::query_flag_truthy(query.briefing.as_deref()) {
return Ok(None);
}
let actor_key = ai_memory_core::ActorKey {
user: actor_user.map(str::to_owned),
session_id: None,
};
Ok(match (handoff_md, brief_md) {
let (ws, proj) = resolve_project_ids_inner(
state,
query.cwd.as_deref(),
query.workspace.as_deref(),
query.project.as_deref(),
ProjectStrategy::parse(query.project_strategy.as_deref()),
&actor_key,
false,
)
.await?;
render_requested_session_brief(state, query, ws, proj).await
}
async fn render_requested_session_brief(
state: &HookState,
query: &HandoffQuery,
workspace_id: WorkspaceId,
project_id: ProjectId,
) -> anyhow::Result<Option<String>> {
if !crate::payload::query_flag_truthy(query.briefing.as_deref()) {
return Ok(None);
}
let budget = query
.briefing_budget
.as_deref()
.and_then(|v| v.trim().parse::<usize>().ok())
.unwrap_or(BRIEF_BUDGET_DEFAULT)
.clamp(BRIEF_BUDGET_MIN, BRIEF_BUDGET_MAX);
let (core, recent) = state
.reader
.session_brief_pages(
workspace_id,
project_id,
BRIEF_CORE_PAGES_LIMIT,
BRIEF_RECENT_PAGES_LIMIT,
)
.await?;
Ok(render_session_brief(&core, &recent, budget))
}
fn combine_handoff_and_brief(
handoff_md: Option<String>,
brief_md: Option<String>,
) -> Option<String> {
match (handoff_md, brief_md) {
(Some(h), Some(b)) => Some(format!("{h}\n{b}")),
(Some(h), None) => Some(h),
(None, Some(b)) => Some(b),
(None, None) => None,
})
}
}
/// Default char budget for the session-start brief (~1k tokens at the
@@ -2263,7 +2310,7 @@ mod tests {
use ai_memory_consolidate::{AutoImproveReviewConfig, run_auto_improve_review};
use ai_memory_core::{SanitizeConfig, Sanitizer};
use ai_memory_llm::{ChatRequest, ChatResponse, LlmProvider, LlmResult};
use ai_memory_store::Store;
use ai_memory_store::{FinishWorkstreamRun, PrepareWorkstreamRun, Store, WorkstreamSelection};
use ai_memory_wiki::Wiki;
use tempfile::TempDir;
@@ -5459,6 +5506,117 @@ mod tests {
);
}
#[tokio::test]
async fn managed_handoff_combines_portable_delta_and_project_brief() {
let tmp = TempDir::new().unwrap();
let state = make_state(&tmp).await;
state
.writer
.upsert_page(brief_page(
state.workspace_id,
state.project_id,
"_rules/managed.md",
"managed briefing sentinel",
false,
))
.await
.unwrap();
let first = state
.writer
.prepare_workstream_run(PrepareWorkstreamRun {
workspace_id: state.workspace_id,
project_id: state.project_id,
repo_fingerprint: "repo".into(),
worktree_fingerprint: "worktree".into(),
cwd: "/repo".into(),
agent: AgentKind::ClaudeCode,
automatic_harness: false,
available_agents: Vec::new(),
selection: WorkstreamSelection::Current,
lease_owner: "test-first".into(),
})
.await
.unwrap();
state
.writer
.finish_workstream_run(FinishWorkstreamRun {
run_id: first.run_id,
native_session_id: Some("claude-session".into()),
source_cursor: None,
events: vec![ai_memory_core::NewWorkstreamEvent {
event_id: "managed-event-1".into(),
agent: AgentKind::ClaudeCode,
native_session_id: "claude-session".into(),
source_record_id: Some("record-1".into()),
kind: WorkstreamEventKind::Message,
role: Some("user".into()),
content: "portable managed delta sentinel".into(),
occurred_at: None,
metadata: serde_json::json!({}),
}],
complete: true,
segment_path: None,
exit_code: Some(0),
})
.await
.unwrap();
let kimi = state
.writer
.prepare_workstream_run(PrepareWorkstreamRun {
workspace_id: state.workspace_id,
project_id: state.project_id,
repo_fingerprint: "repo".into(),
worktree_fingerprint: "worktree".into(),
cwd: "/repo".into(),
agent: AgentKind::KimiCode,
automatic_harness: false,
available_agents: Vec::new(),
selection: WorkstreamSelection::Current,
lease_owner: "test-kimi".into(),
})
.await
.unwrap();
let query = |briefing: Option<&str>| HandoffQuery {
agent: Some("kimi-code".into()),
cwd: Some("/repo".into()),
workspace: Some("default".into()),
project: Some("scratch".into()),
project_strategy: None,
briefing: briefing.map(str::to_owned),
briefing_budget: None,
managed_run: Some(kimi.run_id.to_string()),
session_id: Some("kimi-session".into()),
};
let rendered = fetch_and_accept_handoff(&state, query(Some("true")), None)
.await
.unwrap()
.expect("managed delta and brief must be injected");
let delta_pos = rendered.find("portable managed delta sentinel").unwrap();
let brief_pos = rendered.find("managed briefing sentinel").unwrap();
assert!(
delta_pos < brief_pos,
"managed delta must precede the project brief: {rendered}"
);
let rendered = fetch_and_accept_handoff(&state, query(None), None)
.await
.unwrap();
assert!(
rendered.is_none(),
"delivered managed context must not repeat without a new briefing request"
);
let rendered = fetch_and_accept_handoff(&state, query(Some("true")), None)
.await
.unwrap()
.expect("an explicit later briefing request must still render the project brief");
assert!(rendered.contains("managed briefing sentinel"));
assert!(!rendered.contains("portable managed delta sentinel"));
}
/// The brief renderer respects the char budget: an over-budget body is
/// truncated with a visible note, fully crowded-out core pages are
/// listed as omitted, and an empty project renders nothing at all.
+2 -1
View File
@@ -390,6 +390,7 @@ fn launch_mode(harness: ManagedHarness, args: &[OsString]) -> LaunchMode {
"provider",
"acp",
"web",
"server",
"login",
"doctor",
"vis",
@@ -963,7 +964,7 @@ mod tests {
#[test]
fn kimi_utility_subcommands_are_passed_through() {
for utility in ["export", "doctor", "provider", "upgrade"] {
for utility in ["export", "doctor", "provider", "upgrade", "server"] {
let plan = build_launch_plan(
ManagedHarness::Kimi,
None,
+209 -11
View File
@@ -2,7 +2,7 @@
use std::collections::HashSet;
use std::fs::{self, File};
use std::io::{BufRead as _, BufReader, Seek as _, SeekFrom};
use std::io::{BufRead as _, BufReader, Read as _, Seek as _, SeekFrom};
use std::path::{Path, PathBuf};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
@@ -45,6 +45,11 @@ pub struct ExportedTranscript {
struct FileCursor {
path: String,
offset: u64,
/// Hash of every committed byte through `offset`. Kimi Code can rewrite
/// its journal in place, so its adapter validates this prefix before
/// trusting the byte offset. Other JSONL adapters remain offset-only.
#[serde(default, skip_serializing_if = "Option::is_none")]
prefix_sha256: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
@@ -199,7 +204,23 @@ fn export_jsonl(
let mut file = File::open(path)
.with_context(|| format!("opening native transcript {}", path.display()))?;
let len = file.metadata()?.len();
let start = cursor.map_or(0, |cursor| cursor.offset.min(len));
let (start, mut kimi_prefix_hasher) = if harness == ManagedHarness::Kimi {
let validated = if let Some(cursor) = cursor.as_ref().filter(|cursor| cursor.offset <= len)
&& let Some(expected) = cursor.prefix_sha256.as_deref()
&& let Some(hasher) = hash_file_prefix(&mut file, cursor.offset)?
&& format!("{:x}", hasher.clone().finalize()) == expected
{
Some((cursor.offset, hasher))
} else {
None
};
validated.unwrap_or_else(|| (0, Sha256::new()))
} else {
(
cursor.map_or(0, |cursor| cursor.offset.min(len)),
Sha256::new(),
)
};
file.seek(SeekFrom::Start(start))?;
let mut reader = BufReader::new(file);
let mut offset = start;
@@ -217,6 +238,9 @@ fn export_jsonl(
if !line.ends_with(b"\n") {
break;
}
if harness == ManagedHarness::Kimi {
kimi_prefix_hasher.update(&line);
}
committed_offset = offset;
let value: Value = match serde_json::from_slice(&line) {
Ok(value) => value,
@@ -285,12 +309,31 @@ fn export_jsonl(
source_cursor: Some(serde_json::to_string(&FileCursor {
path: path.to_string_lossy().into_owned(),
offset: committed_offset,
prefix_sha256: (harness == ManagedHarness::Kimi)
.then(|| format!("{:x}", kimi_prefix_hasher.finalize())),
})?),
events,
losses: deduplicate_losses(losses),
})
}
fn hash_file_prefix(file: &mut File, len: u64) -> Result<Option<Sha256>> {
file.seek(SeekFrom::Start(0))?;
let mut hasher = Sha256::new();
let mut remaining = len;
let mut buffer = [0_u8; 16 * 1024];
while remaining > 0 {
let limit = remaining.min(buffer.len() as u64) as usize;
let read = file.read(&mut buffer[..limit])?;
if read == 0 {
return Ok(None);
}
hasher.update(&buffer[..read]);
remaining -= read as u64;
}
Ok(Some(hasher))
}
fn parse_claude(
value: &Value,
session: &str,
@@ -543,12 +586,12 @@ fn parse_pi_family(
}
/// Kimi Code wire journal (`agents/main/wire.jsonl`): flat records
/// `{type, time?, ...payload}`. Only `context.append_message` and
/// `context.apply_compaction` project visible conversation — `turn.prompt`,
/// `turn.steer` and `context.append_loop_event` duplicate the same messages,
/// and records like `config.update`/`llm.request` carry private harness data
/// (system prompts, request bodies) that must never reach the ledger. Unknown
/// record types are ignored so newer kimi versions stay forward-compatible.
/// `{type, time?, ...payload}`. `context.append_message` stores user messages
/// and legacy/imported conversation records; native assistant output and tool
/// exchanges are recorded as `context.append_loop_event`. Records like
/// `config.update`/`llm.request` carry private harness data (system prompts,
/// request bodies) that must never reach the ledger. Unknown record types are
/// ignored so newer Kimi versions stay forward-compatible.
fn parse_kimi(
value: &Value,
session: &str,
@@ -620,12 +663,14 @@ fn parse_kimi(
.and_then(Value::as_array)
.map_or(&[][..], Vec::as_slice);
for (index, call) in tool_calls.iter().enumerate() {
let name = first_string(call, &["name"]).unwrap_or("tool");
let function = call.get("function").unwrap_or(call);
let name = first_string(function, &["name"]).unwrap_or("tool");
// `arguments` is a JSON string; re-serialize it
// compact when it parses, mirroring parse_codex's
// `"{name}: {body}"` tool-call shape.
let arguments = call
.get("arguments")
.or_else(|| function.get("arguments"))
.and_then(Value::as_str)
.map(|raw| match serde_json::from_str::<Value>(raw) {
Ok(parsed) => compact_json(&parsed),
@@ -665,6 +710,84 @@ fn parse_kimi(
_ => {}
}
}
"context.append_loop_event" => {
let event = value.get("event").unwrap_or(&Value::Null);
match event
.get("type")
.and_then(Value::as_str)
.unwrap_or_default()
{
"content.part" => {
let part = event.get("part").unwrap_or(&Value::Null);
match part.get("type").and_then(Value::as_str).unwrap_or_default() {
"text" => {
if let Some(text) = first_string(part, &["text", "content"]) {
push_event(
events,
AgentKind::KimiCode,
session,
record_id,
0,
WorkstreamEventKind::Message,
Some("assistant"),
text,
occurred_at,
json!({}),
);
}
}
"think" | "thinking" => {
losses.push("Kimi hidden reasoning was intentionally excluded".into());
}
_ => {
losses.push(
"Kimi non-text content parts were intentionally excluded".into(),
);
}
}
}
"tool.call" => {
let name = first_string(event, &["name"]).unwrap_or("tool");
let arguments = event.get("args").map(compact_json).unwrap_or_default();
push_event(
events,
AgentKind::KimiCode,
session,
record_id,
0,
WorkstreamEventKind::ToolCall,
Some("assistant"),
&format!("{name}: {arguments}"),
occurred_at,
json!({
"tool": name,
"tool_call_id": event.get("toolCallId").and_then(Value::as_str)
}),
);
}
"tool.result" => {
let result = event.get("result").unwrap_or(&Value::Null);
let texts =
kimi_content_parts(result.get("output").unwrap_or(&Value::Null), losses);
push_event(
events,
AgentKind::KimiCode,
session,
record_id,
0,
WorkstreamEventKind::ToolResult,
Some("tool"),
&texts.parts.join("\n"),
occurred_at,
json!({
"tool_call_id": event.get("toolCallId").and_then(Value::as_str),
"is_error": result.get("isError").and_then(Value::as_bool)
}),
);
}
_ => {}
}
}
"context.apply_compaction" => {
if let Some(summary) = value.get("summary").and_then(Value::as_str) {
push_event(
@@ -726,7 +849,10 @@ struct KimiTextParts {
/// Collect the visible text of a kimi message's content parts in order,
/// annotating parts that cannot be imported (hidden reasoning, media).
fn kimi_text_parts(message: &Value, losses: &mut Vec<String>) -> KimiTextParts {
let content = message.get("content").unwrap_or(&Value::Null);
kimi_content_parts(message.get("content").unwrap_or(&Value::Null), losses)
}
fn kimi_content_parts(content: &Value, losses: &mut Vec<String>) -> KimiTextParts {
let parts: Vec<&Value> = content
.as_array()
.map_or_else(|| vec![content], |items| items.iter().collect());
@@ -2191,10 +2317,14 @@ mod tests {
json!({"type":"context.append_message","message":{"role":"user","origin":{"kind":"hook_result","event":"UserPromptSubmit"},"content":[{"type":"text","text":"injected handoff delta"}]}}),
json!({"type":"context.append_message","message":{"role":"user","origin":{"kind":"injection","variant":"todo"},"content":[{"type":"text","text":"injected todo"}]}}),
json!({"type":"context.append_message","message":{"role":"system","content":[{"type":"text","text":"private system prompt"}]}}),
json!({"type":"context.append_message","message":{"role":"assistant","content":[{"type":"think","think":"private reasoning"},{"type":"text","text":"visible answer"}],"toolCalls":[{"type":"function","id":"call_1","name":"bash","arguments":"{\"cmd\": \"ls\"}"}]}}),
json!({"type":"context.append_message","message":{"role":"assistant","content":[{"type":"think","think":"private reasoning"},{"type":"text","text":"visible answer"}],"toolCalls":[{"type":"function","id":"call_1","function":{"name":"bash","arguments":"{\"cmd\": \"ls\"}"}}]}}),
json!({"type":"context.append_message","message":{"role":"tool","toolCallId":"call_1","content":[{"type":"text","text":"result ok"}]}}),
json!({"type":"context.append_message","message":{"role":"assistant","partial":true,"content":[{"type":"text","text":"stream fragment"}]}}),
json!({"type":"context.apply_compaction","summary":"compact summary","compactedCount":4}),
json!({"type":"context.append_loop_event","event":{"type":"content.part","uuid":"part-think","stepUuid":"step-1","part":{"type":"think","think":"private loop reasoning"}}}),
json!({"type":"context.append_loop_event","event":{"type":"content.part","uuid":"part-text","stepUuid":"step-1","part":{"type":"text","text":"loop visible answer"}}}),
json!({"type":"context.append_loop_event","event":{"type":"tool.call","uuid":"call-2","stepUuid":"step-1","toolCallId":"call_2","name":"Read","args":{"path":"README.md"}}}),
json!({"type":"context.append_loop_event","event":{"type":"tool.result","parentUuid":"call-2","toolCallId":"call_2","result":{"output":[{"type":"text","text":"loop result ok"}],"isError":false}}}),
json!({"type":"config.update","systemPrompt":"never copied"}),
json!({"type":"turn.prompt","text":"duplicate projection"}),
];
@@ -2215,6 +2345,9 @@ mod tests {
WorkstreamEventKind::ToolCall,
WorkstreamEventKind::ToolResult,
WorkstreamEventKind::Compaction,
WorkstreamEventKind::Message,
WorkstreamEventKind::ToolCall,
WorkstreamEventKind::ToolResult,
]
);
assert_eq!(export.events[0].content, "hello kimi");
@@ -2228,6 +2361,12 @@ mod tests {
assert_eq!(export.events[2].metadata["tool"], "bash");
assert_eq!(export.events[3].role.as_deref(), Some("tool"));
assert_eq!(export.events[4].content, "compact summary");
assert_eq!(export.events[5].content, "loop visible answer");
assert_eq!(export.events[6].content, "Read: {\"path\":\"README.md\"}");
assert_eq!(export.events[6].metadata["tool"], "Read");
assert_eq!(export.events[6].metadata["tool_call_id"], "call_2");
assert_eq!(export.events[7].content, "loop result ok");
assert_eq!(export.events[7].metadata["is_error"], false);
assert!(
export
.events
@@ -2281,6 +2420,7 @@ mod tests {
let cursor: FileCursor =
serde_json::from_str(initial.source_cursor.as_deref().unwrap()).unwrap();
assert_eq!(cursor.offset, first.len() as u64 + 1);
assert!(cursor.prefix_sha256.is_some());
fs::write(&wire, format!("{first}\n{second}\n")).unwrap();
let incremental = export_jsonl(
@@ -2294,6 +2434,64 @@ mod tests {
assert_eq!(incremental.events[0].content, "two");
}
#[test]
fn kimi_export_resets_cursor_after_an_in_place_journal_rewrite() {
let temp = tempfile::tempdir().unwrap();
let wire = temp
.path()
.join("store/bucket/session_x/agents/main/wire.jsonl");
fs::create_dir_all(wire.parent().unwrap()).unwrap();
let message = |role: &str, text: &str| {
json!({
"type": "context.append_message",
"message": {
"role": role,
"content": [{"type": "text", "text": text}],
},
})
.to_string()
};
let first = message("user", "one");
let second = message("assistant", "two");
fs::write(&wire, format!("{first}\n{second}\n")).unwrap();
let initial = export_jsonl(ManagedHarness::Kimi, &wire, "session_x", None).unwrap();
let initial_ids: Vec<_> = initial
.events
.iter()
.map(|event| event.event_id.clone())
.collect();
let inserted = message(
"user",
"a rewritten prefix whose different byte length invalidates the old offset",
);
let third = message("assistant", "three");
fs::write(&wire, format!("{inserted}\n{first}\n{second}\n{third}\n")).unwrap();
let rewritten = export_jsonl(
ManagedHarness::Kimi,
&wire,
"session_x",
initial.source_cursor.as_deref(),
)
.unwrap();
assert_eq!(
rewritten
.events
.iter()
.map(|event| event.content.as_str())
.collect::<Vec<_>>(),
[
"a rewritten prefix whose different byte length invalidates the old offset",
"one",
"two",
"three",
]
);
assert_eq!(rewritten.events[1].event_id, initial_ids[0]);
assert_eq!(rewritten.events[2].event_id, initial_ids[1]);
}
#[test]
fn kimi_export_annotates_unimported_subagent_journals() {
let temp = tempfile::tempdir().unwrap();
+5 -3
View File
@@ -665,9 +665,11 @@ compatibility fallback (fire-and-forget POSTs to `/hook`). A pending handoff
is injected at `UserPromptSubmit` through the hook's stdout, which Kimi Code
appends to the model context as a user message before the turn; Kimi Code
fires `SessionStart` but discards that hook's stdout, so hooks installed by
an older release consumed handoffs without delivering them. Re-run
`ai-memory install-hooks --agent kimi-code --apply` after upgrading to pick
up the corrected delivery path (required for `ai-memory run kimi`).
an older release consumed handoffs without delivering them. Existing native
hook commands invoke the current `ai-memory` binary and pick up the corrected
delivery behavior on upgrade. Re-run
`ai-memory install-hooks --agent kimi-code --apply` only for a
script-fallback installation so its staged scripts are refreshed.
Kimi Code hook entries accept only `event`, `matcher`, `command`, and
`timeout`; extra fields make the whole `config.toml` fail to load, so prefer
+10 -4
View File
@@ -175,10 +175,11 @@ ai-memory install-hooks --agent omp --apply
ai-memory install-hooks --agent kimi-code --apply
```
Refreshing the Kimi Code hooks is mandatory when upgrading from a release
before managed kimi support: those hooks fetched the handoff at SessionStart,
whose stdout Kimi Code discards, so pending handoffs were consumed without
ever reaching the model. Current hooks deliver it at UserPromptSubmit.
Kimi Code hooks installed as native `ai-memory hook` commands automatically
pick up the current delivery behavior when the binary is upgraded. A
script-fallback installation must rerun the Kimi Code `install-hooks` command
after upgrading so its staged scripts are refreshed. Current hooks deliver
handoffs at `UserPromptSubmit`; Kimi discards `SessionStart` stdout.
Known Kimi Code adapter limitations: subagent transcripts
(`agents/<id>/wire.jsonl` other than `main`) are not imported in v1 and are
@@ -188,6 +189,11 @@ is a one-way hash of the working directory, so discovery always reads
from the SHA-256 of the raw wire.jsonl line, so two byte-identical lines —
only possible with identical content in the same millisecond, because Kimi
Code stamps each record with `time` — collapse into a single ledger event.
The incremental cursor stores both the complete-record byte offset and a
SHA-256 of that imported prefix. Normal appends resume at the saved offset;
if Kimi rewrites `wire.jsonl` in place, ai-memory resets to the beginning and
replays the file, with stable event ids deduplicating records already in the
workstream.
Legacy sessions that keep `wire.jsonl` directly in the session directory
(the pre-`agents/` layout the kimi session-store still reads through its
stat fallback) are neither discovered nor imported in v1. The native
+3 -2
View File
@@ -88,8 +88,9 @@ default_global = "true"
# brief costs tokens on EVERY session start, so opt in per repo.
# Kimi Code note: kimi discards SessionStart hook stdout, so there the
# brief is delivered on the FIRST user prompt of the session instead
# (once per session, same as Claude); re-briefing after /clear is not
# supported in v1.
# (once per session, same as Claude). Its local delivery markers are created
# only for opted-in repositories and bounded to the 512 newest sessions;
# re-briefing after /clear is not supported in v1.
[briefing]
inject_on_session_start = "true"
+22 -3
View File
@@ -224,8 +224,8 @@ ai_memory_marker_qs() {
# NOT part of ai_memory_marker_qs on purpose: agents that deliver the brief
# once per session (kimi-code, via the first user prompt — kimi discards
# SessionStart hook stdout) append this only on the first fetch, so the
# server does not recompose the brief on every request. Truthiness and the
# char-budget clamp are decided server-side.
# server does not recompose the brief on every request. The char-budget clamp
# is decided server-side.
ai_memory_briefing_qs() {
cwd="$1"
[ -z "$cwd" ] && return 0
@@ -233,8 +233,12 @@ ai_memory_briefing_qs() {
[ -n "$marker" ] || return 0
qs=""
briefing=$(ai_memory_parse_toml_flag "$marker" inject_on_session_start)
case "$(printf '%s' "$briefing" | tr '[:upper:]' '[:lower:]')" in
1|true|yes|on) ;;
*) return 0 ;;
esac
budget=$(ai_memory_parse_toml_flag "$marker" max_chars)
[ -n "$briefing" ] && qs="&briefing=$(ai_memory_url_encode "$briefing")"
qs="&briefing=$(ai_memory_url_encode "$briefing")"
[ -n "$budget" ] && qs="${qs}&briefing_budget=$(ai_memory_url_encode "$budget")"
printf '%s' "$qs"
}
@@ -247,6 +251,21 @@ ai_memory_briefed_file() {
printf '%s/briefed/%s' "$(ai_memory_state_dir)" "$key"
}
# Write a once-per-session briefing marker and keep only the 512 newest
# markers. All marker names are sanitized by ai_memory_briefed_file.
ai_memory_mark_briefed() {
path="$1"
[ -n "$path" ] || return 0
dir=$(dirname "$path")
mkdir -p "$dir" 2>/dev/null || return 0
: > "$path" 2>/dev/null || return 0
LC_ALL=C ls -1t "$dir" 2>/dev/null \
| sed -n '513,$p' \
| while IFS= read -r stale; do
[ -n "$stale" ] && rm -f "$dir/$stale" 2>/dev/null || true
done
}
# Local bridge state for agents whose hook payloads do not carry a session id.
# The value is intentionally non-secret; the server hashes non-UUID ids into its
# typed SessionId domain. `AI_MEMORY_SESSION_ID` may be supplied by advanced
+17 -15
View File
@@ -27,26 +27,28 @@ SESSION_ID=$(ai_memory_extract_session_id "$PAYLOAD")
SESSION_QS=""
[ -n "$SESSION_ID" ] && SESSION_QS="&session_id=$(ai_memory_url_encode "$SESSION_ID")"
# Once-per-session briefing gate, keyed by the native session id (kimi
# always sends `sessionId`); without one, a stable hash of agent+cwd so a
# session-less payload still briefs only once per checkout.
BRIEF_KEY="$SESSION_ID"
if [ -z "$BRIEF_KEY" ]; then
BRIEF_KEY="kimi-code-$(printf '%s' "kimi-code:$CWD" | cksum | awk '{print $1}')"
# Once-per-session briefing gate. Marker files are created only when the
# repository opted in. Prefer Kimi's native session id when supplied;
# otherwise use a stable hash of agent+cwd.
BRIEF_QS=$(ai_memory_briefing_qs "$CWD")
BRIEF_FILE=""
if [ -n "$BRIEF_QS" ]; then
BRIEF_KEY="$SESSION_ID"
if [ -z "$BRIEF_KEY" ]; then
BRIEF_KEY="kimi-code-$(printf '%s' "kimi-code:$CWD" | cksum | awk '{print $1}')"
fi
BRIEF_FILE=$(ai_memory_briefed_file "$BRIEF_KEY")
[ -f "$BRIEF_FILE" ] && BRIEF_QS=""
fi
BRIEF_FILE=$(ai_memory_briefed_file "$BRIEF_KEY")
BRIEF_QS=""
[ -f "$BRIEF_FILE" ] || BRIEF_QS=$(ai_memory_briefing_qs "$CWD")
printf '%s' "$PAYLOAD" \
| ai_memory_post_hook "$SERVER/hook?event=user-prompt&agent=kimi-code${QS}" >/dev/null 2>&1 || true
HANDOFF=$(ai_memory_get_handoff "$SERVER/handoff?agent=kimi-code${QS}${SESSION_QS}${BRIEF_QS}" 2>/dev/null || true)
# Mark the session as briefed only AFTER the GET completed — success or
# error. Fail-open on purpose: with the server down, re-sending the
# brief-flagged request on every prompt would deliver nothing anyway, and
# the one lost brief returns on the next session.
mkdir -p "$(dirname "$BRIEF_FILE")" 2>/dev/null || true
: > "$BRIEF_FILE" 2>/dev/null || true
# Mark an opted-in session as briefed only AFTER the GET completed — success
# or error. Fail-open on purpose: with the server down, re-sending the
# brief-flagged request on every prompt would deliver nothing anyway, and the
# one lost brief returns on the next session.
[ -n "$BRIEF_FILE" ] && ai_memory_mark_briefed "$BRIEF_FILE"
[ -n "$HANDOFF" ] && printf '%s\n' "$HANDOFF"
exit 0
+37 -17
View File
@@ -76,21 +76,27 @@ function Get-AiMemoryTomlFlag {
return $null
}
function Test-AiMemoryTruthy {
param([string] $Value)
if (-not $Value) { return $false }
return @("1", "true", "yes", "on") -contains $Value.Trim().ToLowerInvariant()
}
# Build `&briefing=<v>[&briefing_budget=<v>]` from the `[briefing]` section
# of the marker walked up from $Cwd. Returns "" when the repo did not opt
# in. Used by agents that deliver the compiled project brief once per
# session (kimi-code, via the first user prompt — kimi discards
# SessionStart hook stdout) so the server does not recompose the brief on
# every request. Truthiness and the char-budget clamp are server-side.
# every request. The char-budget clamp is server-side.
function Get-AiMemoryBriefingQuery {
param([string] $Cwd)
if (-not $Cwd) { return "" }
$marker = Get-AiMemoryMarkerToml -Cwd $Cwd
if (-not $marker) { return "" }
$qs = ""
$briefing = Get-AiMemoryTomlFlag -File $marker -Key "inject_on_session_start"
if (-not (Test-AiMemoryTruthy -Value $briefing)) { return "" }
$budget = Get-AiMemoryTomlFlag -File $marker -Key "max_chars"
if ($briefing) { $qs += "&briefing=$([uri]::EscapeDataString($briefing))" }
$qs = "&briefing=$([uri]::EscapeDataString($briefing))"
if ($budget) { $qs += "&briefing_budget=$([uri]::EscapeDataString($budget))" }
return $qs
}
@@ -104,6 +110,18 @@ function Get-AiMemoryBriefedFile {
return (Join-Path (Join-Path (Get-AiMemoryStateDir) "briefed") $safe)
}
function Set-AiMemoryBriefed {
param([string] $Path)
if (-not $Path) { return }
$dir = Split-Path $Path -Parent
New-Item -ItemType Directory -Force -Path $dir -ErrorAction SilentlyContinue | Out-Null
New-Item -ItemType File -Force -Path $Path -ErrorAction SilentlyContinue | Out-Null
Get-ChildItem -Path $dir -File -ErrorAction SilentlyContinue |
Sort-Object LastWriteTimeUtc -Descending |
Select-Object -Skip 512 |
Remove-Item -Force -ErrorAction SilentlyContinue
}
# Resolve the basename of the MAIN git repository root for $Cwd, following the
# worktree commondir pointer so every linked worktree collapses to one stable
# name. Mirrors the POSIX `ai_memory_repo_root_project`: a containerized server
@@ -290,21 +308,24 @@ function Invoke-AiMemoryHook {
}
} catch {
}
# Once-per-session briefing gate, keyed by the native session id
# (kimi always sends `sessionId`); without one, a stable hash of
# agent+cwd so a session-less payload still briefs only once.
# Once-per-session briefing gate. Marker files are created only for
# repositories that opt in. Prefer the native session id when Kimi
# supplies one; otherwise use a stable hash of agent+cwd.
$BriefQS = ""
$BriefFile = $null
if ($BriefingOncePerSession) {
$BriefKey = [string]$NativeSessionId
if (-not $BriefKey) {
$Sha = [System.Security.Cryptography.SHA256]::Create()
$Bytes = $Sha.ComputeHash([System.Text.Encoding]::UTF8.GetBytes("$Agent`n$Cwd"))
$BriefKey = (($Bytes | ForEach-Object { $_.ToString("x2") }) -join "")
}
$BriefFile = Get-AiMemoryBriefedFile -Key $BriefKey
if (-not (Test-Path $BriefFile -PathType Leaf)) {
$BriefQS = Get-AiMemoryBriefingQuery -Cwd $Cwd
$BriefQS = Get-AiMemoryBriefingQuery -Cwd $Cwd
if ($BriefQS) {
$BriefKey = [string]$NativeSessionId
if (-not $BriefKey) {
$Sha = [System.Security.Cryptography.SHA256]::Create()
$Bytes = $Sha.ComputeHash([System.Text.Encoding]::UTF8.GetBytes("$Agent`n$Cwd"))
$BriefKey = (($Bytes | ForEach-Object { $_.ToString("x2") }) -join "")
}
$BriefFile = Get-AiMemoryBriefedFile -Key $BriefKey
if (Test-Path $BriefFile -PathType Leaf) {
$BriefQS = ""
}
}
}
try {
@@ -335,8 +356,7 @@ function Invoke-AiMemoryHook {
# brief-flagged request on every prompt would deliver nothing
# anyway, and the one lost brief returns on the next session).
if ($BriefFile) {
New-Item -ItemType Directory -Force -Path (Split-Path $BriefFile -Parent) -ErrorAction SilentlyContinue | Out-Null
New-Item -ItemType File -Force -Path $BriefFile -ErrorAction SilentlyContinue | Out-Null
Set-AiMemoryBriefed -Path $BriefFile
}
} elseif ($AntigravityPreInvocationOutput) {
[Console]::Out.Write("{}")
+9 -2
View File
@@ -550,11 +550,18 @@ done
# Kimi Code keeps providers/model and hooks in one config.toml under
# $KIMI_CODE_HOME. Seed the isolated home with the operator's provider
# settings; install-hooks merges its [[hooks]] entries without rewriting
# the rest of the file.
# settings and minimum login state; install-hooks merges its [[hooks]]
# entries without rewriting the rest of the file. Native sessions, indexes,
# logs, and telemetry remain isolated.
if [ -f "$HOME/.kimi-code/config.toml" ]; then
cp "$HOME/.kimi-code/config.toml" "$KIMI_ACCEPTANCE_HOME/config.toml"
fi
for relative in credentials/kimi-code.json oauth/kimi-code device_id; do
if [ -f "$HOME/.kimi-code/$relative" ]; then
mkdir -p "$KIMI_ACCEPTANCE_HOME/$(dirname "$relative")"
cp "$HOME/.kimi-code/$relative" "$KIMI_ACCEPTANCE_HOME/$relative"
fi
done
install_hook() {
local agent=$1