use std::collections::HashMap;
use std::io::{BufRead, BufReader, Read, Write};
use std::process::{Child, ChildStdin, Command, Stdio};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::mpsc::{self, Receiver, Sender};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use anyhow::{Context, Result, anyhow, bail};
use serde_json::{Value, json};
use super::transcript;
use crate::agent::Agent;
use crate::context::{self, CLOSE, OPEN};
use crate::explain::{DEEP_ASK_ABOVE_TOKENS, Progress, Request};
const SYSTEM: &str = "You are \"peek\", an explainer embedded in a terminal. The user highlighted a fragment \
of text (marked ⟦like this⟧) in their GitHub Copilot CLI session and wants to understand it. Explain \
what the fragment means where it appears: define the terms, name the specific thing it refers to when \
the context shows it, say what any code does there, and why it matters for the user's task. No \
preamble, no headings. Answer from the provided context and general knowledge only.";
fn copilot_bin() -> String {
std::env::var("PEEKME_COPILOT_BIN").unwrap_or_else(|_| {
crate::launch::find_real("copilot")
.map(|p| p.to_string_lossy().into_owned())
.unwrap_or_else(|| "copilot".into())
})
}
fn model() -> Option<String> {
std::env::var("PEEKME_COPILOT_MODEL")
.ok()
.filter(|m| !m.trim().is_empty())
}
type Pending = Mutex<HashMap<u64, Sender<Result<Value, String>>>>;
pub struct Server {
stdin: Mutex<ChildStdin>,
child: Mutex<Child>,
next_id: AtomicU64,
pending: Arc<Pending>,
subscribers: Arc<Mutex<HashMap<String, Sender<Value>>>>,
}
impl Server {
fn start() -> Result<Arc<Self>> {
let mut child = Command::new(copilot_bin())
.args([
"--headless",
"--stdio",
"--no-auto-update",
"--log-level",
"none",
])
.current_dir(std::env::temp_dir())
.env_remove(crate::launch::ACTIVE_ENV)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::null())
.spawn()
.context("could not start `copilot --headless`")?;
let stdin = child.stdin.take().unwrap();
let stdout = child.stdout.take().unwrap();
let server = Arc::new(Self {
stdin: Mutex::new(stdin),
child: Mutex::new(child),
next_id: AtomicU64::new(1),
pending: Arc::default(),
subscribers: Arc::default(),
});
let reader = server.clone();
std::thread::spawn(move || reader.read_loop(BufReader::new(stdout)));
server.request("ping", json!({"message": "peekme"}))?;
Ok(server)
}
fn read_loop(&self, mut out: impl BufRead) {
while let Some(msg) = read_frame(&mut out) {
match (msg.get("id"), msg.get("method")) {
(Some(id), None) => {
let Some(id) = id.as_u64() else { continue };
if let Some(tx) = self.pending.lock().unwrap().remove(&id) {
let res = match msg.get("error") {
Some(e) => Err(e
.get("message")
.and_then(Value::as_str)
.unwrap_or("error")
.to_string()),
None => Ok(msg.get("result").cloned().unwrap_or(Value::Null)),
};
let _ = tx.send(res);
}
}
(Some(id), Some(_)) => {
let _ = self.send(&json!({"jsonrpc": "2.0", "id": id,
"error": {"code": -32601, "message": "not supported by peekme"}}));
}
(None, Some(_)) => {
let sid = msg.pointer("/params/sessionId").and_then(Value::as_str);
if let Some(sid) = sid
&& let Some(tx) = self.subscribers.lock().unwrap().get(sid)
{
let _ = tx.send(msg);
}
}
_ => {}
}
}
for (_, tx) in self.pending.lock().unwrap().drain() {
let _ = tx.send(Err("the Copilot server exited".into()));
}
self.subscribers.lock().unwrap().clear();
}
fn send(&self, msg: &Value) -> Result<()> {
let body = msg.to_string();
let mut stdin = self.stdin.lock().unwrap();
write!(stdin, "Content-Length: {}\r\n\r\n{body}", body.len())?;
stdin.flush()?;
Ok(())
}
fn request(&self, method: &str, params: Value) -> Result<Value> {
let id = self.next_id.fetch_add(1, Ordering::Relaxed);
let (tx, rx) = mpsc::channel();
self.pending.lock().unwrap().insert(id, tx);
self.send(&json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params}))?;
rx.recv_timeout(Duration::from_secs(60))
.map_err(|_| anyhow!("Copilot did not answer `{method}`"))?
.map_err(|e| anyhow!("{e}"))
}
fn subscribe(&self, sid: &str) -> Receiver<Value> {
let (tx, rx) = mpsc::channel();
self.subscribers.lock().unwrap().insert(sid.to_string(), tx);
rx
}
fn unsubscribe(&self, sid: &str) {
self.subscribers.lock().unwrap().remove(sid);
}
fn is_alive(&self) -> bool {
matches!(self.child.lock().unwrap().try_wait(), Ok(None))
}
fn shutdown(&self) {
let mut child = self.child.lock().unwrap();
let _ = child.kill();
let _ = child.wait();
}
}
fn read_frame(out: &mut impl BufRead) -> Option<Value> {
loop {
let mut len = None;
loop {
let mut line = String::new();
if out.read_line(&mut line).ok()? == 0 {
return None;
}
let line = line.trim();
if line.is_empty() {
break;
}
if let Some((k, v)) = line.split_once(':')
&& k.eq_ignore_ascii_case("content-length")
{
len = v.trim().parse::<usize>().ok();
}
}
let Some(len) = len else { continue };
let mut body = vec![0u8; len];
out.read_exact(&mut body).ok()?;
if let Ok(v) = serde_json::from_slice(&body) {
return Some(v);
}
}
}
#[derive(Default)]
pub struct Slot(Mutex<Option<Arc<Server>>>);
impl Slot {
fn get(&self) -> Result<Arc<Server>> {
let mut guard = self.0.lock().unwrap();
if let Some(s) = guard.as_ref().filter(|s| s.is_alive()) {
return Ok(s.clone());
}
if let Some(old) = guard.take() {
old.shutdown();
}
let server = Server::start()?;
*guard = Some(server.clone());
Ok(server)
}
pub fn shutdown(&self) {
if let Some(s) = self.0.lock().unwrap().take() {
s.shutdown();
}
}
}
pub fn prompt(req: &Request) -> (context::Built, Option<usize>) {
let pids = req
.agent_pid
.map(crate::launch::process_tree)
.unwrap_or_default();
let conv = transcript::config_dir()
.and_then(|c| transcript::find(&c, &pids, &req.cwd))
.and_then(|p| transcript::load(&p));
let hit = conv.as_ref().and_then(|c| context::find(c, &req.screen));
let compact = context::build(Agent::Copilot, &req.cwd, conv.as_ref(), hit, &req.screen);
let (mut built, compact_len) = match (&conv, req.deep) {
(Some(c), true) => (
context::build_deep(Agent::Copilot, &req.cwd, c, hit, &req.screen),
Some(compact.prompt.chars().count()),
),
_ => (compact, None),
};
built.prompt.push_str(&format!(
"Answer in at most {} short lines; keep the {OPEN}{CLOSE} marks out of the answer.\n",
req.max_lines
));
(built, compact_len)
}
pub fn explain(slot: &Slot, req: Request, tx: impl Fn(Progress)) {
if let Err(e) = run(slot, req, &tx) {
tx(Progress::Failed(format!("{e:#}")));
}
}
fn run(slot: &Slot, req: Request, tx: &impl Fn(Progress)) -> Result<()> {
let (built, compact_len) = prompt(&req);
if let Some(compact) = compact_len
&& !req.force
{
let (deep, normal) = (built.prompt.chars().count() / 4 + 300, compact / 4 + 300);
if deep > DEEP_ASK_ABOVE_TOKENS {
tx(Progress::TooBig {
tokens: deep,
ratio: deep / normal,
});
return Ok(());
}
}
if std::env::var_os("PEEKME_DEBUG_PROMPT").is_some() {
let _ = std::fs::write(
std::env::temp_dir().join("peekme-last-prompt.txt"),
&built.prompt,
);
}
let server = slot.get()?;
let auth = server.request("auth.getStatus", json!({}))?;
if auth.get("isAuthenticated").and_then(Value::as_bool) != Some(true) {
bail!("Copilot CLI is not signed in. Run `copilot login`.");
}
let sid = uuid_v4();
let events = server.subscribe(&sid);
let result = (|| -> Result<()> {
let mut params = json!({
"sessionId": sid,
"clientName": "peekme",
"availableTools": [],
"streaming": true,
"reasoningEffort": "none",
"systemMessage": {"mode": "replace", "content": SYSTEM},
});
if let Some(m) = model() {
params["model"] = json!(m);
}
server.request("session.create", params)?;
server.request(
"session.send",
json!({"sessionId": sid, "prompt": built.prompt}),
)?;
let mut model_name = model().unwrap_or_else(|| "Copilot".into());
let mut started = false;
loop {
let msg = events
.recv_timeout(Duration::from_secs(90))
.map_err(|_| anyhow!("the explanation timed out"))?;
let event = msg.pointer("/params/event").cloned().unwrap_or(Value::Null);
let data = event.get("data").cloned().unwrap_or(Value::Null);
match event.get("type").and_then(Value::as_str) {
Some("session.auto_mode_resolved") => {
if let Some(m) = data.get("chosenModel").and_then(Value::as_str) {
model_name = m.to_string();
}
}
Some("assistant.message_delta") => {
if let Some(d) = data.get("deltaContent").and_then(Value::as_str) {
if !started {
started = true;
tx(Progress::Started {
model: model_name.clone(),
source: built.label,
});
}
tx(Progress::Delta(d.to_string()));
}
}
Some("assistant.message") if !started => {
if let Some(t) = data.get("content").and_then(Value::as_str) {
started = true;
tx(Progress::Started {
model: model_name.clone(),
source: built.label,
});
tx(Progress::Delta(t.to_string()));
}
}
Some("session.error") => {
let m = data
.get("message")
.and_then(Value::as_str)
.unwrap_or("Copilot returned an error");
bail!("{m}");
}
Some("session.idle") => {
tx(Progress::Done);
return Ok(());
}
_ => {}
}
}
})();
server.unsubscribe(&sid);
let _ = server.request("session.delete", json!({"sessionId": sid}));
result
}
fn uuid_v4() -> String {
let mut b = [0u8; 16];
if let Ok(mut f) = std::fs::File::open("/dev/urandom") {
let _ = f.read_exact(&mut b);
}
if b == [0u8; 16] {
let t = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_nanos());
b = (t ^ (u128::from(std::process::id()) << 64)).to_le_bytes();
}
b[6] = (b[6] & 0x0f) | 0x40;
b[8] = (b[8] & 0x3f) | 0x80;
let h: String = b.iter().map(|x| format!("{x:02x}")).collect();
format!(
"{}-{}-{}-{}-{}",
&h[0..8],
&h[8..12],
&h[12..16],
&h[16..20],
&h[20..32]
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn frames_and_uuids() {
let bytes = b"Content-Length: 14\r\n\r\n{\"id\":1,\"a\":2}Content-Type: x\r\ncontent-length: 2\r\n\r\n{}";
let mut r = BufReader::new(&bytes[..]);
assert_eq!(read_frame(&mut r), Some(json!({"id": 1, "a": 2})));
assert_eq!(read_frame(&mut r), Some(json!({})));
assert_eq!(read_frame(&mut r), None);
let u = uuid_v4();
assert_eq!(u.len(), 36);
assert_eq!(&u[14..15], "4");
assert_ne!(uuid_v4(), u);
}
}