The PRD's first acceptance gate now holds: grep -RinE '\bx\.ai\b|grok' crates/ --include='*.rs' → 0 matches (exempt: NOTICE and third-party license archives, README provenance, and the required 'Based on Grok Build Open Source' attribution, now sourced from version_attribution.txt). Wire-visible renames (both sides in this repo, changed in lockstep): - Auth method id 'grok.com' → 'kimi-code' (AuthMethodKind::KimiCode). - Every x.ai/* and _x.ai/* ACP ext method and meta key → kigi/* / _kigi/* (~200 names; grokShell → kigiShell). Session-file replay keeps a read-side alias for the legacy '_x.ai/session/update' method so existing updates.jsonl histories load; writes emit only the new name (both directions test-pinned). - Agent types grok-build* → kigi* with a documented legacy-prefix alias at resolution time so persisted sessions keep resolving. - ToolNamespace/BuiltinAgentName GrokBuild* → Kigi* (wire snake_case kigi/kigi_concise/kigi_hashline; schema regenerated); grok_build implementation dirs renamed to kigi*. - x-grok-* headers → x-kigi-*, __GROK_* sentinels → __KIGI_*, themes grokday/groknight → kigiday/kiginight (old persisted values fall back to the default theme), web_fetch allowlist xAI hosts → kimi.com + moonshot platforms, changelog CDN → this repo, grok-build changelog archives deleted. - BYOK default endpoint removed: [endpoints] api_base_url is now truly optional with NO default — consumers fail fast with the flag name when unset (no silent x.ai egress). Mock harnesses inject it explicitly. - System-prompt identity fixed: 'released by xAI' → 'an unofficial community CLI for Kimi' (template + regenerated encrypted form). Also repaired pre-existing grok-era test debt found by the sweep: the stale trace_classify default-model pin, the grok-pager UA label test, pty-harness stale-binary reuse and non-hermetic moonshot routing (a PTY test could previously reach the real api.moonshot.cn), and the outdated oauth fixture scope key. Gates: §9 grep 0; fmt clean; workspace check/clippy 0/0 (-D warnings); FULL cargo test --workspace: 234 suites, 21,961 passed, 0 failed; deny advisories ok.
142 lines
4.4 KiB
Rust
142 lines
4.4 KiB
Rust
//! Unified log forwarding for the pager.
|
|
//!
|
|
//! Buffers log entries in memory and flushes them to the shell via
|
|
//! `kigi/log` ACP notifications. Call [`init`] once at startup with
|
|
//! the ACP sender, then use [`info`], [`warn`], [`error`], [`debug`]
|
|
//! from anywhere.
|
|
|
|
use std::sync::{Mutex, OnceLock};
|
|
|
|
use agent_client_protocol as acp;
|
|
use kigi_acp_lib::AcpAgentTx;
|
|
use kigi_log::unified_log::{
|
|
ClientLogEntry, LOG_METHOD, LogLevel, LogNotificationParams, LogSource,
|
|
};
|
|
use tokio::runtime::Handle;
|
|
|
|
static ACP_TX: OnceLock<AcpAgentTx> = OnceLock::new();
|
|
static BUFFER: Mutex<Vec<ClientLogEntry>> = Mutex::new(Vec::new());
|
|
|
|
/// Initialize the unified log forwarder with the ACP sender.
|
|
///
|
|
/// Must be called once after the ACP connection is established.
|
|
/// Spawns a background task that flushes buffered entries every few
|
|
/// seconds so events are delivered promptly without manual flush calls.
|
|
/// Entries buffered before this call will be picked up on the first tick.
|
|
pub fn init(tx: AcpAgentTx) {
|
|
let _ = ACP_TX.set(tx);
|
|
tokio::spawn(async {
|
|
let mut interval = tokio::time::interval(std::time::Duration::from_secs(1));
|
|
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
|
|
loop {
|
|
interval.tick().await;
|
|
flush();
|
|
}
|
|
});
|
|
}
|
|
|
|
fn now_ts() -> String {
|
|
chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Millis, true)
|
|
}
|
|
|
|
fn push_entry(lvl: LogLevel, msg: &str, sid: Option<&str>, ctx: Option<serde_json::Value>) {
|
|
let entry = ClientLogEntry {
|
|
ts: now_ts(),
|
|
pid: Some(std::process::id()),
|
|
ver: Some(kigi_version::VERSION.to_owned()),
|
|
lvl,
|
|
sid: sid.map(Into::into),
|
|
msg: msg.into(),
|
|
ctx,
|
|
};
|
|
if let Ok(mut buf) = BUFFER.lock() {
|
|
buf.push(entry);
|
|
// Auto-flush when we have a reasonable batch, but only if the ACP
|
|
// sender is ready -- otherwise keep buffering until an explicit flush().
|
|
if buf.len() >= 16 && ACP_TX.get().is_some() {
|
|
let entries: Vec<ClientLogEntry> = buf.drain(..).collect();
|
|
drop(buf);
|
|
send_entries(entries);
|
|
}
|
|
}
|
|
}
|
|
|
|
fn build_notification(entries: Vec<ClientLogEntry>) -> Option<acp::ExtNotification> {
|
|
if entries.is_empty() {
|
|
return None;
|
|
}
|
|
let params = LogNotificationParams {
|
|
src: LogSource::KigiPager,
|
|
entries,
|
|
};
|
|
let raw = serde_json::value::to_raw_value(¶ms).ok()?;
|
|
Some(acp::ExtNotification::new(LOG_METHOD, raw.into()))
|
|
}
|
|
|
|
fn send_entries(entries: Vec<ClientLogEntry>) {
|
|
let Some(tx) = ACP_TX.get() else { return };
|
|
let Some(notification) = build_notification(entries) else {
|
|
return;
|
|
};
|
|
// Guard against panic if called from a non-tokio thread (e.g., a
|
|
// tracing::Layer callback on a blocking thread).
|
|
let Ok(handle) = Handle::try_current() else {
|
|
return;
|
|
};
|
|
let tx = tx.clone();
|
|
handle.spawn(async move {
|
|
let _ = kigi_acp_lib::acp_send(notification, &tx).await;
|
|
});
|
|
}
|
|
|
|
/// Flush any buffered entries to the shell (fire-and-forget).
|
|
pub fn flush() {
|
|
let entries = {
|
|
let Ok(mut buf) = BUFFER.lock() else { return };
|
|
if buf.is_empty() {
|
|
return;
|
|
}
|
|
buf.drain(..).collect::<Vec<_>>()
|
|
};
|
|
send_entries(entries);
|
|
}
|
|
|
|
/// Flush buffered entries and await delivery.
|
|
///
|
|
/// Use this before process exit to ensure entries are delivered
|
|
/// before the agent shuts down.
|
|
pub async fn flush_blocking() {
|
|
let entries = {
|
|
let Ok(mut buf) = BUFFER.lock() else { return };
|
|
if buf.is_empty() {
|
|
return;
|
|
}
|
|
buf.drain(..).collect::<Vec<_>>()
|
|
};
|
|
let Some(tx) = ACP_TX.get() else { return };
|
|
let Some(notification) = build_notification(entries) else {
|
|
return;
|
|
};
|
|
let _ = kigi_acp_lib::acp_send(notification, tx).await;
|
|
}
|
|
|
|
/// Log an info-level entry.
|
|
pub fn info(msg: &str, sid: Option<&str>, ctx: Option<serde_json::Value>) {
|
|
push_entry(LogLevel::Info, msg, sid, ctx);
|
|
}
|
|
|
|
/// Log a warn-level entry.
|
|
pub fn warn(msg: &str, sid: Option<&str>, ctx: Option<serde_json::Value>) {
|
|
push_entry(LogLevel::Warn, msg, sid, ctx);
|
|
}
|
|
|
|
/// Log an error-level entry.
|
|
pub fn error(msg: &str, sid: Option<&str>, ctx: Option<serde_json::Value>) {
|
|
push_entry(LogLevel::Error, msg, sid, ctx);
|
|
}
|
|
|
|
/// Log a debug-level entry.
|
|
pub fn debug(msg: &str, sid: Option<&str>, ctx: Option<serde_json::Value>) {
|
|
push_entry(LogLevel::Debug, msg, sid, ctx);
|
|
}
|