mirror of
https://github.com/l0ng-ai/tty7.git
synced 2026-09-22 00:02:23 +00:00
Merge remote-tracking branch 'origin/main' into feat/ssh-settings-redesign
This commit is contained in:
@@ -5,6 +5,25 @@ All notable changes to tty7 are documented in this file.
|
||||
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),
|
||||
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Added
|
||||
|
||||
- **Pi is a first-class agent, not a fallback one** — Pi panes drew the generic
|
||||
robot glyph every unbranded agent shares, so a Pi tab was indistinguishable
|
||||
from an Aider or Qwen one in the sidebar, the tab chip and the tray menu. They
|
||||
now carry their own avatar on the existing sky accent, status dot unchanged.
|
||||
The mark is Pi's own, from pi.dev, rescaled to tty7's 24x24 icon grid — its
|
||||
`prefers-color-scheme` stylesheet dropped, since these avatars are tinted by
|
||||
the app. Restoring a Pi pane also resumes its conversation now: the tty7
|
||||
extension reports Pi's session id, and the resume command is
|
||||
`pi --session <id>` (Pi's `--resume` is a boolean that only opens the
|
||||
interactive picker), with `--session` / `--session-id` / `--fork` /
|
||||
`--resume` / `-r` / `--continue` / `-c` stripped off the replayed launch
|
||||
flags so the restored id wins. A pane launched with `--no-session` is not
|
||||
resumed at all — that pane never wrote a session to disk, and reopening one
|
||||
would override the choice to keep it ephemeral. (#225)
|
||||
|
||||
## [26.7.6] - 2026-07-28
|
||||
|
||||
### Added
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
<!-- Pi's own brand mark, from https://pi.dev/logo-auto.svg.
|
||||
Normalized to tty7's 24x24 icon format; geometry unchanged. The published
|
||||
file is an 800x800 box carrying a `prefers-color-scheme` <style> block that
|
||||
swaps black for white — no other mark in this set ships theme-conditional
|
||||
CSS, and the avatars are tinted by the app regardless, so the class and
|
||||
media query are dropped for the flat sentinel fill the rest of the set
|
||||
uses. The mark lands on an exact 4x4 grid of 6-unit cells at this size. -->
|
||||
<svg width="24" height="24" viewBox="0 0 24 24" fill="none" xmlns="http://www.w3.org/2000/svg">
|
||||
<path fill-rule="evenodd" clip-rule="evenodd" d="M0 0H18V12H12V18H6V24H0V0ZM6 6V12H12V6H6Z" fill="#FF0000"/>
|
||||
<path d="M18 12H24V24H18V12Z" fill="#FF0000"/>
|
||||
</svg>
|
||||
|
After Width: | Height: | Size: 795 B |
@@ -1201,7 +1201,17 @@ export const Tty7Presence = async ({{ $ }}) => {{
|
||||
/// extensions from per-directory `index.ts` files; this one forwards Pi's
|
||||
/// lifecycle events to the emitter. Inert outside tty7 (both the TS guard and
|
||||
/// the emitter check `TTY7`).
|
||||
///
|
||||
/// Every forwarded event carries Pi's own session id, read off the read-only
|
||||
/// session manager on the handler's context — that id is what
|
||||
/// [`crate::core::cli_agent::CLIAgent::resume_command`] feeds to `pi --session
|
||||
/// <id>` when a restored pane relaunches its conversation. The context is only
|
||||
/// reachable from a handler, so the load-time presence ping stays bare and
|
||||
/// `session_start` (fired for startup / new / resume / fork) is what actually
|
||||
/// reports the id, including when the user switches sessions mid-pane.
|
||||
fn pi_extension_ts(target: &HookTarget) -> Option<String> {
|
||||
// The machine that will *run* the agent, not this one — a remote workspace
|
||||
// installs its hooks over there, where `current_exe` is the wrong binary.
|
||||
let exe = serde_json::to_string(&target.exe.display().to_string()).ok()?;
|
||||
Some(format!(
|
||||
r#"/* tty7 agent-hook pi bridge — generated by tty7, do not edit. */
|
||||
@@ -1210,20 +1220,42 @@ import {{ spawnSync }} from "node:child_process";
|
||||
|
||||
const EXE = {exe};
|
||||
|
||||
function emit(event: string): void {{
|
||||
/** The slice of Pi's handler context we read — structural, so this bridge does
|
||||
* not depend on the context type staying exported. */
|
||||
type SessionCtx = {{ sessionManager?: {{ getSessionId?(): string | undefined }} }};
|
||||
|
||||
function emit(event: string, ctx?: SessionCtx): void {{
|
||||
try {{
|
||||
spawnSync(EXE, ["agent-hook", "pi", event], {{ stdio: ["ignore", "ignore", "ignore"] }});
|
||||
let payload = "";
|
||||
try {{
|
||||
const id = ctx?.sessionManager?.getSessionId?.();
|
||||
if (id) payload = JSON.stringify({{ session_id: id }});
|
||||
}} catch {{}}
|
||||
const args = ["agent-hook", "pi", event];
|
||||
// Nothing to send → leave stdin closed rather than handing the emitter a
|
||||
// pipe it has to read to EOF.
|
||||
if (payload) {{
|
||||
spawnSync(EXE, args, {{ input: payload, stdio: ["pipe", "ignore", "ignore"] }});
|
||||
}} else {{
|
||||
spawnSync(EXE, args, {{ stdio: ["ignore", "ignore", "ignore"] }});
|
||||
}}
|
||||
}} catch {{}}
|
||||
}}
|
||||
|
||||
export default function (pi: ExtensionAPI) {{
|
||||
if (!process.env["TTY7"]) return;
|
||||
// Extension load = the agent is running in this pane; Pi has no separate
|
||||
// session-start event.
|
||||
// Extension load = the agent is running in this pane. No context here yet,
|
||||
// so the id rides on session_start instead.
|
||||
emit("session-start");
|
||||
pi.on("agent_start", () => emit("prompt-submit"));
|
||||
pi.on("agent_end", () => emit("stop"));
|
||||
pi.on("session_shutdown", () => emit("session-end"));
|
||||
pi.on("agent_start", (_event, ctx) => emit("prompt-submit", ctx));
|
||||
pi.on("agent_end", (_event, ctx) => emit("stop", ctx));
|
||||
pi.on("session_shutdown", (_event, ctx) => emit("session-end", ctx));
|
||||
// Last, and guarded: the three above already worked, so a Pi build that
|
||||
// rejects this event name must not take them — or the whole extension —
|
||||
// down with it.
|
||||
try {{
|
||||
pi.on("session_start", (_event, ctx) => emit("session-start", ctx));
|
||||
}} catch {{}}
|
||||
}}
|
||||
"#
|
||||
))
|
||||
@@ -1632,6 +1664,14 @@ mod tests {
|
||||
assert!(pi.contains("agent-hook pi"));
|
||||
assert!(pi.contains(&exe));
|
||||
assert!(pi.contains(r#"process.env["TTY7"]"#));
|
||||
// Pi's session id is what `pi --session <id>` resumes with, and it only
|
||||
// reaches tty7 if the bridge pipes a payload instead of ignoring stdin
|
||||
// — which is how this integration shipped, silently costing Pi panes
|
||||
// their resume.
|
||||
assert!(pi.contains("getSessionId"));
|
||||
assert!(pi.contains("session_id"));
|
||||
assert!(pi.contains(r#"stdio: ["pipe", "ignore", "ignore"]"#));
|
||||
assert!(pi.contains(r#"pi.on("session_start""#));
|
||||
|
||||
let grok = grok_hooks_json(&target).expect("grok content builds");
|
||||
let parsed: serde_json::Value = serde_json::from_str(&grok).expect("valid JSON");
|
||||
|
||||
@@ -183,6 +183,11 @@ impl CLIAgent {
|
||||
{
|
||||
return None;
|
||||
}
|
||||
// A pane the user launched as deliberately ephemeral has nothing on
|
||||
// disk to come back to, whatever id the agent reported.
|
||||
if launch_argv.is_some_and(|argv| self.opts_out_of_sessions(argv)) {
|
||||
return None;
|
||||
}
|
||||
// The user's launch flags, pre-joined with a leading space so they
|
||||
// splice into the format strings below; empty when none survive.
|
||||
let flags = launch_argv
|
||||
@@ -214,10 +219,32 @@ impl CLIAgent {
|
||||
// Grok Build: `grok --resume <id-or-title>`; a UUID-shaped value
|
||||
// always takes the id path, which is what its hooks report.
|
||||
CLIAgent::Grok => Some(format!("grok{flags} --resume {session_id}")),
|
||||
// Pi's `--resume`/`-r` is a *boolean* that opens the interactive
|
||||
// session picker and `--continue`/`-c` just takes the newest
|
||||
// session; the flag that targets one by id is `--session
|
||||
// <path|id>` ("Use specific session file or partial UUID"). Its
|
||||
// ids are uuidv7, so they clear the token gate above.
|
||||
CLIAgent::Pi => Some(format!("pi{flags} --session {session_id}")),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Whether `argv` launched the agent with session persistence turned off,
|
||||
/// which makes the pane unresumable: nothing was written to disk, so a
|
||||
/// replayed id would point at a session file that never existed *and*
|
||||
/// would quietly undo the user's opt-out. Distinct from the stale flags in
|
||||
/// [`Self::replay_flags`], which name a different session and merely have
|
||||
/// to lose to the injected id.
|
||||
fn opts_out_of_sessions(self, argv: &[String]) -> bool {
|
||||
let ephemeral: &[&str] = match self {
|
||||
// Pi still mints an in-memory session id under `--no-session` — it
|
||||
// only skips the write — so tty7 does observe an id to replay.
|
||||
CLIAgent::Pi => &["--no-session"],
|
||||
_ => &[],
|
||||
};
|
||||
argv.iter().any(|t| ephemeral.contains(&t.as_str()))
|
||||
}
|
||||
|
||||
/// The launch-flag tail of `argv` worth replaying on a resume command, or
|
||||
/// `None` to resume bare. Deliberately conservative: anything ambiguous
|
||||
/// falls back to no flags rather than a corrupted command line.
|
||||
@@ -271,6 +298,23 @@ impl CLIAgent {
|
||||
// `--last` targets "the most recent session" and would contradict
|
||||
// the explicit id we inject.
|
||||
CLIAgent::Codex => &["--last"],
|
||||
// Pi's ways of picking a session all fight the `--session <id>` we
|
||||
// inject: `--session`/`--session-id` name a different one,
|
||||
// `--fork` would branch instead of continue, and the boolean
|
||||
// `-r`/`-c` re-open the picker or the newest session.
|
||||
// `--session-dir` is *not* here — it says where sessions live, so
|
||||
// the id we inject needs it to still be there. `--no-session` is
|
||||
// not here either: it isn't stale, it means there is nothing to
|
||||
// resume at all (see [`Self::opts_out_of_sessions`]).
|
||||
CLIAgent::Pi => &[
|
||||
"--session",
|
||||
"--session-id",
|
||||
"--fork",
|
||||
"--resume",
|
||||
"-r",
|
||||
"--continue",
|
||||
"-c",
|
||||
],
|
||||
// Beyond the session-targeting flags (`--load` is grok's hidden
|
||||
// alias for `--resume`; `--session-id` names a *new* session and
|
||||
// `--fork-session` would branch off the one we mean to continue),
|
||||
@@ -383,9 +427,9 @@ impl CLIAgent {
|
||||
CLIAgent::Goose => "icons/agents/goose.svg",
|
||||
CLIAgent::Droid => "icons/agents/droid.svg",
|
||||
CLIAgent::Grok => "icons/agents/grok.svg",
|
||||
CLIAgent::Pi => "icons/agents/pi.svg",
|
||||
// No brand mark bundled → generic robot glyph.
|
||||
CLIAgent::Aider
|
||||
| CLIAgent::Pi
|
||||
| CLIAgent::Auggie
|
||||
| CLIAgent::Hermes
|
||||
| CLIAgent::Vibe
|
||||
@@ -951,6 +995,32 @@ mod tests {
|
||||
assert_eq!(CLIAgent::Grok.accent_rgb(), 0x000000);
|
||||
}
|
||||
|
||||
/// The fallback robot glyph is a placeholder, not a resting state: an agent
|
||||
/// may only sit on it while no mark is bundled, and a bundled mark may not
|
||||
/// silently fall back off. Pinning the exact fallback set makes either
|
||||
/// direction a deliberate edit here rather than something noticed in the UI.
|
||||
#[test]
|
||||
fn only_the_unbranded_agents_use_the_fallback_glyph() {
|
||||
let fallback: Vec<&str> = CLIAgent::ALL
|
||||
.into_iter()
|
||||
.filter(|a| a.icon_path() == "icons/bot.svg")
|
||||
.map(CLIAgent::slug)
|
||||
.collect();
|
||||
assert_eq!(
|
||||
fallback,
|
||||
["aider", "auggie", "hermes", "vibe", "antigravity", "qwen"]
|
||||
);
|
||||
// Everything else names a bundled mark under the agents directory.
|
||||
for a in CLIAgent::ALL {
|
||||
let path = a.icon_path();
|
||||
assert!(
|
||||
path == "icons/bot.svg" || path == format!("icons/agents/{}.svg", a.slug()),
|
||||
"{} points at an unexpected {path}",
|
||||
a.display_name()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn detects_newer_agents_by_command() {
|
||||
for (cmd, agent) in [
|
||||
@@ -1241,6 +1311,14 @@ mod tests {
|
||||
CLIAgent::Codex.resume_command("th_read.9", None).as_deref(),
|
||||
Some("codex resume th_read.9")
|
||||
);
|
||||
// Pi targets a session by id through `--session`, not through its
|
||||
// boolean `--resume` (which only opens the picker).
|
||||
assert_eq!(
|
||||
CLIAgent::Pi
|
||||
.resume_command("0199c3f2-1b0e-7c3a-9f21-6d4b8e2a5c17", None)
|
||||
.as_deref(),
|
||||
Some("pi --session 0199c3f2-1b0e-7c3a-9f21-6d4b8e2a5c17")
|
||||
);
|
||||
// No resume flag known → None.
|
||||
assert_eq!(CLIAgent::Aider.resume_command("abc", None), None);
|
||||
// An id carrying shell syntax is refused outright.
|
||||
@@ -1362,6 +1440,60 @@ mod tests {
|
||||
.as_deref(),
|
||||
Some("codex resume id-3 --yolo")
|
||||
);
|
||||
// Pi: mode flags replay, but every way of naming a *different* session
|
||||
// is stripped so the injected id wins — including the boolean picker
|
||||
// flags.
|
||||
assert_eq!(
|
||||
CLIAgent::Pi
|
||||
.resume_command("id-a", Some(&argv(&["pi", "--model", "opus"])))
|
||||
.as_deref(),
|
||||
Some("pi --model opus --session id-a")
|
||||
);
|
||||
assert_eq!(
|
||||
CLIAgent::Pi
|
||||
.resume_command(
|
||||
"id-b",
|
||||
Some(&argv(&[
|
||||
"pi",
|
||||
"--session",
|
||||
"old-id",
|
||||
"--fork",
|
||||
"old",
|
||||
"-c",
|
||||
"--model",
|
||||
"opus"
|
||||
]))
|
||||
)
|
||||
.as_deref(),
|
||||
Some("pi --model opus --session id-b")
|
||||
);
|
||||
// `--no-session` is the user asking for an ephemeral pane: Pi mints an
|
||||
// id but never writes the file, so resuming it would open an empty
|
||||
// session *and* override the opt-out — no resume command at all.
|
||||
assert_eq!(
|
||||
CLIAgent::Pi.resume_command(
|
||||
"id-x",
|
||||
Some(&argv(&["pi", "--no-session", "--model", "opus"]))
|
||||
),
|
||||
None
|
||||
);
|
||||
// `--session-dir` says where sessions live — the injected id needs it,
|
||||
// so it is deliberately *not* stripped.
|
||||
assert_eq!(
|
||||
CLIAgent::Pi
|
||||
.resume_command(
|
||||
"id-c",
|
||||
Some(&argv(&[
|
||||
"pi",
|
||||
"--session-dir",
|
||||
"/w/.sessions",
|
||||
"--fork",
|
||||
"old"
|
||||
]))
|
||||
)
|
||||
.as_deref(),
|
||||
Some("pi --session-dir /w/.sessions --session id-c")
|
||||
);
|
||||
// No token names the agent (custom wrapper rule) → bare.
|
||||
assert_eq!(
|
||||
CLIAgent::Claude
|
||||
|
||||
@@ -72,6 +72,20 @@ pub const MAX_RECORD_BYTES: usize = 4 * 1024 * 1024;
|
||||
/// error rather than grow the file until the disk fills.
|
||||
pub const MAX_WORKSPACES: usize = 1024;
|
||||
|
||||
/// Ceiling on the whole document, which is what a single `WorkspaceList` reply
|
||||
/// has to fit into.
|
||||
///
|
||||
/// The per-record and per-count ceilings above are independent of each other,
|
||||
/// and their product is 4 GiB — sixty-four times the frame limit. Seventeen
|
||||
/// accepted `WorkspacePut`s of a maximal record are enough to put the array
|
||||
/// past it, and from then on *every* `WorkspaceList` on the machine is a reply
|
||||
/// that cannot be encoded: every client shows an empty workspace list, and the
|
||||
/// only repair is editing the file by hand. So the total is bounded where it is
|
||||
/// actually known — at the save — with room to spare under
|
||||
/// [`MAX_FRAME`](crate::daemon::protocol::MAX_FRAME), since what is measured
|
||||
/// here is the pretty-printed form and the wire carries the compact one.
|
||||
pub const MAX_STORE_BYTES: usize = 32 * 1024 * 1024;
|
||||
|
||||
/// Ceiling on a record key, which is a workspace uuid in every non-hostile
|
||||
/// case.
|
||||
const MAX_ID_BYTES: usize = 128;
|
||||
@@ -172,6 +186,23 @@ struct State {
|
||||
/// because the file's array order is what a client lists, and a hash map
|
||||
/// would reshuffle the picker on every save for no reason.
|
||||
records: Vec<(String, Value)>,
|
||||
/// `(mtime, len)` of the file as this snapshot last saw it, or `None` when
|
||||
/// there was no file.
|
||||
///
|
||||
/// This store is not always the only writer. The design's answer is one
|
||||
/// server per machine, and `tty7-server --stdio` now starts the daemon
|
||||
/// rather than serving in-process for exactly that reason — but an explicit
|
||||
/// `--serve`, or a daemon that could not be started, still leaves two
|
||||
/// processes over one file. `persist` writes the *whole* document, so
|
||||
/// without noticing that the file moved underneath it, the second to save
|
||||
/// silently drops everything the first did.
|
||||
stamp: Option<(std::time::SystemTime, u64)>,
|
||||
}
|
||||
|
||||
/// The file's identity as far as [`State::stamp`] is concerned.
|
||||
fn stamp_of(path: &Path) -> Option<(std::time::SystemTime, u64)> {
|
||||
let meta = std::fs::metadata(path).ok()?;
|
||||
Some((meta.modified().ok()?, meta.len()))
|
||||
}
|
||||
|
||||
impl WorkspaceStore {
|
||||
@@ -186,9 +217,10 @@ impl WorkspaceStore {
|
||||
pub fn open(path: impl Into<PathBuf>) -> Arc<WorkspaceStore> {
|
||||
let path = path.into();
|
||||
let records = load_records(&path);
|
||||
let stamp = stamp_of(&path);
|
||||
Arc::new(WorkspaceStore {
|
||||
state: Mutex::new(State { records, stamp }),
|
||||
path,
|
||||
state: Mutex::new(State { records }),
|
||||
attachments: Mutex::new(Vec::new()),
|
||||
subscribers: Mutex::new(Vec::new()),
|
||||
next_subscriber: AtomicU64::new(1),
|
||||
@@ -302,7 +334,7 @@ impl WorkspaceStore {
|
||||
Undo::Remove(st.records.len() - 1)
|
||||
}
|
||||
};
|
||||
if let Err(e) = self.persist(&st) {
|
||||
if let Err(e) = self.persist(&st, true) {
|
||||
match undo {
|
||||
Undo::Restore(i, old) => st.records[i].1 = old,
|
||||
Undo::Remove(i) => {
|
||||
@@ -311,6 +343,7 @@ impl WorkspaceStore {
|
||||
}
|
||||
return Err(e);
|
||||
}
|
||||
self.restamp(&mut st);
|
||||
}
|
||||
|
||||
self.notify(id, origin);
|
||||
@@ -329,10 +362,11 @@ impl WorkspaceStore {
|
||||
return Ok(false);
|
||||
};
|
||||
let removed = st.records.remove(i);
|
||||
if let Err(e) = self.persist(&st) {
|
||||
if let Err(e) = self.persist(&st, false) {
|
||||
st.records.insert(i, removed);
|
||||
return Err(e);
|
||||
}
|
||||
self.restamp(&mut st);
|
||||
}
|
||||
// The attachment goes with it: nothing can be attached to a workspace
|
||||
// that no longer exists, and leaving the entry would have M6 report a
|
||||
@@ -440,7 +474,24 @@ impl WorkspaceStore {
|
||||
// in-memory state is still a valid state (the undo path restores it
|
||||
// before returning) and the file is either the old or the new one, so
|
||||
// carrying on is strictly better than taking the server down.
|
||||
self.state.lock().unwrap_or_else(|e| e.into_inner())
|
||||
let mut st = self.state.lock().unwrap_or_else(|e| e.into_inner());
|
||||
|
||||
// Re-read when the file moved under us. Cheap — one `stat` — and it is
|
||||
// what keeps a second writer's changes from being overwritten by this
|
||||
// store's whole-document save, since the base we mutate is then theirs
|
||||
// rather than a snapshot from before their write. It also lets a read
|
||||
// see their changes at all: `notify` reaches subscribers in *this*
|
||||
// process only.
|
||||
let on_disk = stamp_of(&self.path);
|
||||
if on_disk != st.stamp {
|
||||
log::debug!(
|
||||
"{} changed underneath this store; re-reading",
|
||||
self.path.display()
|
||||
);
|
||||
st.records = load_records(&self.path);
|
||||
st.stamp = on_disk;
|
||||
}
|
||||
st
|
||||
}
|
||||
|
||||
fn attachments_locked(&self) -> std::sync::MutexGuard<'_, Vec<(String, Attachment)>> {
|
||||
@@ -453,7 +504,14 @@ impl WorkspaceStore {
|
||||
/// [`Workspaces`](crate::core::session::Workspaces) parses, so this file is
|
||||
/// readable by the same code that reads a client's `session.json` and a
|
||||
/// human can diff the two.
|
||||
fn persist(&self, st: &State) -> io::Result<()> {
|
||||
/// Write the whole document.
|
||||
///
|
||||
/// `bounded` asks for [`MAX_STORE_BYTES`] to be enforced. Set by the paths
|
||||
/// that *grow* the file and clear by the ones that shrink it: a store that
|
||||
/// came up holding an over-large file — written by an older build, or by
|
||||
/// hand — must still be able to delete its way back under the limit rather
|
||||
/// than refusing every operation including the repair.
|
||||
fn persist(&self, st: &State, bounded: bool) -> io::Result<()> {
|
||||
#[derive(Serialize)]
|
||||
struct Doc<'a> {
|
||||
workspaces: Vec<&'a Value>,
|
||||
@@ -462,11 +520,28 @@ impl WorkspaceStore {
|
||||
workspaces: st.records.iter().map(|(_, v)| v).collect(),
|
||||
};
|
||||
let bytes = serde_json::to_vec_pretty(&doc).map_err(io::Error::other)?;
|
||||
if bounded && bytes.len() > MAX_STORE_BYTES {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::InvalidInput,
|
||||
format!(
|
||||
"the workspace store would be {} bytes; the limit is {MAX_STORE_BYTES}, \
|
||||
which is what one WorkspaceList reply has to fit into",
|
||||
bytes.len()
|
||||
),
|
||||
));
|
||||
}
|
||||
if let Some(parent) = self.path.parent() {
|
||||
std::fs::create_dir_all(parent)?;
|
||||
}
|
||||
crate::core::config::write_atomic(&self.path, &bytes)
|
||||
}
|
||||
|
||||
/// Record the file's identity after this store wrote it, so the next
|
||||
/// [`WorkspaceStore::locked`] does not mistake its own save for someone
|
||||
/// else's and re-read it.
|
||||
fn restamp(&self, st: &mut State) {
|
||||
st.stamp = stamp_of(&self.path);
|
||||
}
|
||||
}
|
||||
|
||||
/// How to undo a mutation whose write failed.
|
||||
@@ -632,6 +707,102 @@ mod tests {
|
||||
})
|
||||
}
|
||||
|
||||
/// Two stores over one file — an explicit `--serve` alongside a daemon, or
|
||||
/// a daemon that could not be started — must not silently undo each other.
|
||||
///
|
||||
/// `persist` writes the whole document, so a store that mutates a snapshot
|
||||
/// taken before the other's write puts that stale snapshot back. This is
|
||||
/// how a workspace rename made on the laptop vanishes the next time the
|
||||
/// desktop reorders a tab, with nothing reported to either.
|
||||
#[test]
|
||||
fn a_second_writer_does_not_get_overwritten_by_a_stale_snapshot() {
|
||||
let dir = tempfile::TempDir::new().unwrap();
|
||||
let path = dir.path().join(STORE_FILE);
|
||||
let first = WorkspaceStore::open(&path);
|
||||
let second = WorkspaceStore::open(&path);
|
||||
|
||||
first.put("w1", record("w1", "one"), None).unwrap();
|
||||
first.put("w2", record("w2", "two"), None).unwrap();
|
||||
|
||||
// `second` last read the file when it was empty. It has to notice.
|
||||
second
|
||||
.put("w2", record("w2", "two, renamed"), None)
|
||||
.unwrap();
|
||||
|
||||
let names: Vec<String> = WorkspaceStore::open(&path)
|
||||
.list()
|
||||
.iter()
|
||||
.map(|r| r["name"].as_str().unwrap_or_default().to_string())
|
||||
.collect();
|
||||
assert_eq!(
|
||||
names,
|
||||
["one", "two, renamed"],
|
||||
"the second writer's save dropped what the first had written"
|
||||
);
|
||||
}
|
||||
|
||||
/// The same, one layer down: a read sees another process's write, because
|
||||
/// `notify` only ever reaches subscribers inside this process.
|
||||
#[test]
|
||||
fn a_read_sees_a_change_another_store_made_to_the_file() {
|
||||
let dir = tempfile::TempDir::new().unwrap();
|
||||
let path = dir.path().join(STORE_FILE);
|
||||
let reader = WorkspaceStore::open(&path);
|
||||
let writer = WorkspaceStore::open(&path);
|
||||
|
||||
assert!(reader.get("w1").is_none());
|
||||
writer.put("w1", record("w1", "one"), None).unwrap();
|
||||
assert_eq!(
|
||||
reader.get("w1").map(|r| r["name"].clone()),
|
||||
Some(serde_json::json!("one")),
|
||||
"a read answered from a snapshot older than the file"
|
||||
);
|
||||
}
|
||||
|
||||
/// The per-record and per-count ceilings do not bound their product, so the
|
||||
/// document is bounded where it is known — at the save.
|
||||
///
|
||||
/// Past `MAX_FRAME` the store is not merely large, it is unreadable: every
|
||||
/// `WorkspaceList` becomes a reply that cannot be encoded, so every client
|
||||
/// shows an empty list and the only repair is editing the file by hand.
|
||||
#[test]
|
||||
fn a_put_that_would_outgrow_one_reply_is_refused_and_undone() {
|
||||
let (store, _dir) = store();
|
||||
// Records big enough that a handful crosses the limit, and small enough
|
||||
// that the test stays quick.
|
||||
let chunk = "x".repeat(2 * 1024 * 1024);
|
||||
let big = |id: &str| {
|
||||
let mut r = record(id, "big");
|
||||
r["padding"] = Value::String(chunk.clone());
|
||||
r
|
||||
};
|
||||
|
||||
let mut accepted = 0;
|
||||
let refusal = loop {
|
||||
let id = format!("w{accepted}");
|
||||
match store.put(&id, big(&id), None) {
|
||||
Ok(()) => accepted += 1,
|
||||
Err(e) => break e,
|
||||
}
|
||||
assert!(accepted < 64, "the total was never bounded");
|
||||
};
|
||||
assert_eq!(refusal.kind(), io::ErrorKind::InvalidInput);
|
||||
assert!(
|
||||
refusal.to_string().contains("WorkspaceList"),
|
||||
"the refusal has to say what the limit is for: {refusal}"
|
||||
);
|
||||
|
||||
// Refused, not half-applied: the record that did not fit is not in the
|
||||
// store and is not in the file.
|
||||
assert_eq!(store.len(), accepted);
|
||||
assert!(store.get(&format!("w{accepted}")).is_none());
|
||||
assert_eq!(WorkspaceStore::open(store.path()).len(), accepted);
|
||||
|
||||
// And a delete still works, so a store that came up over the limit can
|
||||
// be repaired rather than being wedged.
|
||||
assert!(store.delete("w0", None).unwrap());
|
||||
}
|
||||
|
||||
// ── The basics ──────────────────────────────────────────────────────────
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -76,7 +76,7 @@ use std::time::{Duration, Instant};
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use super::protocol::{read_frame, write_frame};
|
||||
use super::protocol::{MAX_FRAME, read_frame, write_frame};
|
||||
|
||||
/// The control dialect's own version, negotiated in [`ControlHello`] and
|
||||
/// independent of [`crate::daemon::protocol::PROTOCOL_VERSION`]: the pane
|
||||
@@ -777,24 +777,49 @@ pub enum ControlClientMsg {
|
||||
impl ControlClientMsg {
|
||||
/// Encode and write this message as one frame.
|
||||
pub fn encode<W: Write>(&self, w: &mut W) -> io::Result<()> {
|
||||
match self {
|
||||
ControlClientMsg::Hello(hello) => write_frame(w, kind::HELLO, &to_json(hello)?),
|
||||
let (k, payload) = self.to_frame()?;
|
||||
write_frame(w, k, &payload)
|
||||
}
|
||||
|
||||
/// The frame this message *would* write, without writing it.
|
||||
///
|
||||
/// Split out so a caller can tell "this message cannot be encoded" from
|
||||
/// "the link failed". Everything fallible here happens with the wire
|
||||
/// untouched, and the difference matters: marking the connection dead over
|
||||
/// a local serialization failure puts the workspace into `Reconnecting`
|
||||
/// over a request the server never saw.
|
||||
pub fn to_frame(&self) -> io::Result<(u8, Vec<u8>)> {
|
||||
let (k, payload) = match self {
|
||||
ControlClientMsg::Hello(hello) => (kind::HELLO, to_json(hello)?),
|
||||
ControlClientMsg::Request { req_id, req } => {
|
||||
require_nonzero(*req_id, "CONTROL_REQUEST")?;
|
||||
let body = encode_body(*req_id, &to_json(req)?, &[])?;
|
||||
write_frame(w, kind::REQUEST, &body)
|
||||
(kind::REQUEST, encode_body(*req_id, &to_json(req)?, &[])?)
|
||||
}
|
||||
ControlClientMsg::RequestBlob { req_id, req, blob } => {
|
||||
require_nonzero(*req_id, "CONTROL_REQUEST_BLOB")?;
|
||||
let body = encode_body(*req_id, &to_json(req)?, blob)?;
|
||||
write_frame(w, kind::REQUEST_BLOB, &body)
|
||||
(
|
||||
kind::REQUEST_BLOB,
|
||||
encode_body(*req_id, &to_json(req)?, blob)?,
|
||||
)
|
||||
}
|
||||
ControlClientMsg::Cancel { req_id } => {
|
||||
require_nonzero(*req_id, "CONTROL_CANCEL")?;
|
||||
let body = encode_body(*req_id, &[], &[])?;
|
||||
write_frame(w, kind::CANCEL, &body)
|
||||
(kind::CANCEL, encode_body(*req_id, &[], &[])?)
|
||||
}
|
||||
};
|
||||
// See `ControlServerMsg::to_frame`: reaching the verdict here rather
|
||||
// than inside `write_frame` is what keeps "too big" a failure with the
|
||||
// wire untouched.
|
||||
if payload.len() > MAX_FRAME {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::InvalidData,
|
||||
format!(
|
||||
"control frame of {} bytes exceeds the {MAX_FRAME}-byte limit",
|
||||
payload.len()
|
||||
),
|
||||
));
|
||||
}
|
||||
Ok((k, payload))
|
||||
}
|
||||
|
||||
/// Decode one already-read frame.
|
||||
@@ -858,12 +883,24 @@ pub enum ControlServerMsg {
|
||||
impl ControlServerMsg {
|
||||
/// Encode and write this message as one frame.
|
||||
pub fn encode<W: Write>(&self, w: &mut W) -> io::Result<()> {
|
||||
match self {
|
||||
ControlServerMsg::HelloOk(ok) => write_frame(w, kind::HELLO_OK, &to_json(ok)?),
|
||||
let (k, payload) = self.to_frame()?;
|
||||
write_frame(w, k, &payload)
|
||||
}
|
||||
|
||||
/// The frame this message *would* write, without writing it.
|
||||
///
|
||||
/// The server's half of the same split as
|
||||
/// [`ControlClientMsg::to_frame`], and it earns its keep in the same way:
|
||||
/// a reply that cannot be encoded — a `SearchHit` path that is not UTF-8, a
|
||||
/// payload past [`MAX_FRAME`] — leaves the wire untouched, so the server can
|
||||
/// still answer the request with the error instead of dropping the reply and
|
||||
/// leaving the client to wait out its deadline.
|
||||
pub fn to_frame(&self) -> io::Result<(u8, Vec<u8>)> {
|
||||
let (k, payload) = match self {
|
||||
ControlServerMsg::HelloOk(ok) => (kind::HELLO_OK, to_json(ok)?),
|
||||
ControlServerMsg::Response { req_id, reply } => {
|
||||
require_nonzero(*req_id, "CONTROL_RESPONSE")?;
|
||||
let body = encode_body(*req_id, &to_json(reply)?, &[])?;
|
||||
write_frame(w, kind::RESPONSE, &body)
|
||||
(kind::RESPONSE, encode_body(*req_id, &to_json(reply)?, &[])?)
|
||||
}
|
||||
ControlServerMsg::ResponseBlob {
|
||||
req_id,
|
||||
@@ -871,14 +908,26 @@ impl ControlServerMsg {
|
||||
blob,
|
||||
} => {
|
||||
require_nonzero(*req_id, "CONTROL_RESPONSE_BLOB")?;
|
||||
let body = encode_body(*req_id, &to_json(reply)?, blob)?;
|
||||
write_frame(w, kind::RESPONSE_BLOB, &body)
|
||||
}
|
||||
ControlServerMsg::Event(event) => {
|
||||
let body = encode_body(0, &to_json(event)?, &[])?;
|
||||
write_frame(w, kind::EVENT, &body)
|
||||
(
|
||||
kind::RESPONSE_BLOB,
|
||||
encode_body(*req_id, &to_json(reply)?, blob)?,
|
||||
)
|
||||
}
|
||||
ControlServerMsg::Event(event) => (kind::EVENT, encode_body(0, &to_json(event)?, &[])?),
|
||||
};
|
||||
// Checked here rather than left to `write_frame`, so that "too big to
|
||||
// send" is a verdict reached with nothing written — which is what lets
|
||||
// the caller answer with an error instead of going silent.
|
||||
if payload.len() > MAX_FRAME {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::InvalidData,
|
||||
format!(
|
||||
"control frame of {} bytes exceeds the {MAX_FRAME}-byte limit",
|
||||
payload.len()
|
||||
),
|
||||
));
|
||||
}
|
||||
Ok((k, payload))
|
||||
}
|
||||
|
||||
/// Decode one already-read frame.
|
||||
@@ -1348,11 +1397,19 @@ impl ClientInner {
|
||||
}
|
||||
|
||||
fn send(&self, msg: &ControlClientMsg) -> io::Result<()> {
|
||||
// Encoded before the writer is even locked, and before anything is
|
||||
// concluded about the link. A message that cannot be serialized, or a
|
||||
// payload past `MAX_FRAME`, fails with the wire untouched — the server
|
||||
// never saw the request, so the connection is exactly as healthy as it
|
||||
// was and must not be marked dead. Only a failure *after* bytes could
|
||||
// have gone out leaves the stream in a state worth giving up on.
|
||||
let (k, payload) = msg.to_frame()?;
|
||||
|
||||
let mut w = self
|
||||
.writer
|
||||
.lock()
|
||||
.map_err(|_| io::Error::other("control writer was poisoned"))?;
|
||||
let r = msg.encode(&mut *w).and_then(|()| w.flush());
|
||||
let r = write_frame(&mut *w, k, &payload).and_then(|()| w.flush());
|
||||
if r.is_err() {
|
||||
self.connected.store(false, Ordering::Release);
|
||||
}
|
||||
@@ -1387,12 +1444,17 @@ impl ClientInner {
|
||||
log::trace!("control reply for unknown req_id {req_id}; dropping");
|
||||
return;
|
||||
};
|
||||
drop(pending);
|
||||
// Under the `pending` lock, which is the one `forget` takes first.
|
||||
// Inserting after releasing it loses the race with a caller that times
|
||||
// out in the gap: its `forget` clears a `blobs` entry that does not
|
||||
// exist yet, this insert lands with nobody left to take it, and a whole
|
||||
// file's contents stay in the map for the life of the connection.
|
||||
if !blob.is_empty()
|
||||
&& let Ok(mut blobs) = self.blobs.lock()
|
||||
{
|
||||
blobs.insert(req_id, blob);
|
||||
}
|
||||
drop(pending);
|
||||
// Capacity 1 and one reply per id, so this cannot block; it can only
|
||||
// fail if the caller already gave up between the table lookup and here.
|
||||
let _ = tx.try_send(reply);
|
||||
@@ -1487,6 +1549,143 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
/// A bare `ClientInner` with no link behind it — enough for the side-table
|
||||
/// bookkeeping, which is all `deliver`/`forget` touch.
|
||||
fn detached_inner() -> Arc<ClientInner> {
|
||||
Arc::new(ClientInner {
|
||||
writer: Mutex::new(Box::new(Cursor::new(Vec::new()))),
|
||||
next_req_id: AtomicU64::new(1),
|
||||
pending: Mutex::new(HashMap::new()),
|
||||
blobs: Mutex::new(HashMap::new()),
|
||||
connected: AtomicBool::new(true),
|
||||
last_inbound: Mutex::new(Instant::now()),
|
||||
hello: ControlHelloOk {
|
||||
control_version: CONTROL_VERSION,
|
||||
protocol_version: 0,
|
||||
build: "test".into(),
|
||||
separator: '/',
|
||||
home: "/home/me".into(),
|
||||
features: Vec::new(),
|
||||
},
|
||||
shutdown: None,
|
||||
reader_done: Mutex::new(false),
|
||||
reader_exit: Condvar::new(),
|
||||
})
|
||||
}
|
||||
|
||||
/// A path its own platform accepts but `str` cannot hold — Latin-1 bytes on
|
||||
/// Unix, an unpaired surrogate on Windows. Both are real filenames, and
|
||||
/// neither survives a JSON wire.
|
||||
fn unrepresentable_path() -> PathBuf {
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use std::os::unix::ffi::OsStrExt as _;
|
||||
PathBuf::from(std::ffi::OsStr::from_bytes(b"/tmp/caf\xe9.rs"))
|
||||
}
|
||||
#[cfg(windows)]
|
||||
{
|
||||
use std::os::windows::ffi::OsStringExt as _;
|
||||
// `C:\` then a lone high surrogate: legal in an NTFS name, not
|
||||
// convertible to a `str`.
|
||||
PathBuf::from(std::ffi::OsString::from_wide(&[
|
||||
0x0043, 0x003a, 0x005c, 0xd800, 0x002e, 0x0072, 0x0073,
|
||||
]))
|
||||
}
|
||||
}
|
||||
|
||||
/// `to_frame` reaches every "cannot be sent" verdict with the wire
|
||||
/// untouched. That is what lets the server answer such a reply with an
|
||||
/// error instead of dropping it, and the client keep its link on a local
|
||||
/// encode failure — both of which are silent hangs otherwise.
|
||||
#[test]
|
||||
fn to_frame_refuses_what_cannot_be_sent_without_writing_it() {
|
||||
// A path the platform accepts and `str` cannot hold. `serde` errors on
|
||||
// such a `Path` rather than converting it lossily, which is the shape a
|
||||
// `SearchHit` takes when the server has one under the search roots.
|
||||
let unrepresentable = ControlServerMsg::Response {
|
||||
req_id: 1,
|
||||
reply: ControlReply::Ok(ReplyOk::Hits(vec![SearchHit {
|
||||
name: "caf\u{fffd}.rs".into(),
|
||||
path: unrepresentable_path(),
|
||||
is_dir: false,
|
||||
ignored: false,
|
||||
}])),
|
||||
};
|
||||
assert!(
|
||||
unrepresentable.to_frame().is_err(),
|
||||
"a path that is not valid Unicode must not encode"
|
||||
);
|
||||
|
||||
// And the size verdict, which `write_frame` would otherwise reach only
|
||||
// after the writer was locked.
|
||||
let huge = ControlServerMsg::ResponseBlob {
|
||||
req_id: 1,
|
||||
reply: ControlReply::Ok(ReplyOk::Meta(meta())),
|
||||
blob: vec![0u8; MAX_FRAME + 1],
|
||||
};
|
||||
let err = huge
|
||||
.to_frame()
|
||||
.expect_err("an oversize reply must not encode");
|
||||
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
|
||||
}
|
||||
|
||||
/// A message that cannot be put on the wire is not a dead link. `send`
|
||||
/// fails before the writer is even locked, so the server never saw the
|
||||
/// request and the connection is exactly as healthy as it was — marking it
|
||||
/// dead would drop the workspace into `Reconnecting` and, since
|
||||
/// `is_connected` never returns to true, keep it there.
|
||||
#[test]
|
||||
fn a_frame_too_large_to_send_does_not_condemn_the_link() {
|
||||
let inner = detached_inner();
|
||||
let oversize = ControlClientMsg::RequestBlob {
|
||||
req_id: 1,
|
||||
req: ControlRequest::WriteFile {
|
||||
path: "/home/me/huge.bin".into(),
|
||||
},
|
||||
blob: vec![0u8; MAX_FRAME + 1],
|
||||
};
|
||||
|
||||
let err = inner
|
||||
.send(&oversize)
|
||||
.expect_err("an oversize frame must fail");
|
||||
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
|
||||
assert!(
|
||||
inner.connected.load(Ordering::Acquire),
|
||||
"a local encode failure marked the connection dead"
|
||||
);
|
||||
|
||||
// And the link still works for a message that does fit.
|
||||
inner.send(&ControlClientMsg::Cancel { req_id: 1 }).unwrap();
|
||||
}
|
||||
|
||||
/// A caller giving up at the same instant its reply lands must not leave
|
||||
/// the blob behind. `deliver` and `forget` both touch two tables, and if
|
||||
/// the blob goes in after the `pending` lock is released, `forget` clears
|
||||
/// an entry that does not exist yet and the bytes — a whole file, up to the
|
||||
/// editor's limit — are retained for the life of the connection.
|
||||
///
|
||||
/// Hammered rather than interleaved by hand: there is no seam to inject at,
|
||||
/// so this leans on repetition to hit the window. It fails intermittently
|
||||
/// against the insert-after-unlock ordering and never against the fix.
|
||||
#[test]
|
||||
fn a_blob_racing_its_callers_timeout_is_not_retained() {
|
||||
for req_id in 1..2_000u64 {
|
||||
let inner = detached_inner();
|
||||
let (tx, _rx) = std::sync::mpsc::sync_channel(1);
|
||||
inner.pending.lock().unwrap().insert(req_id, tx);
|
||||
|
||||
let giving_up = Arc::clone(&inner);
|
||||
let t = std::thread::spawn(move || giving_up.forget(req_id));
|
||||
inner.deliver(req_id, ControlReply::Ok(ReplyOk::Unit), vec![7u8; 64]);
|
||||
t.join().unwrap();
|
||||
|
||||
assert!(
|
||||
inner.blobs.lock().unwrap().is_empty(),
|
||||
"a blob outlived the request it belonged to (req_id {req_id})"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
fn every_request() -> Vec<ControlRequest> {
|
||||
vec![
|
||||
ControlRequest::Ping,
|
||||
|
||||
@@ -600,15 +600,20 @@ mod tests {
|
||||
|
||||
/// Pre-epoch times round-trip exactly rather than clamping to zero, because
|
||||
/// the editor compares mtimes for equality.
|
||||
///
|
||||
/// Every nanosecond figure here is a multiple of 100: a Windows
|
||||
/// `SystemTime` is a FILETIME, whose tick *is* 100ns, so a finer value
|
||||
/// would be rounded on the way in and the assertion would be about
|
||||
/// `SystemTime`'s resolution rather than about this conversion.
|
||||
#[test]
|
||||
fn mtime_handles_both_sides_of_the_epoch() {
|
||||
use std::time::{Duration, UNIX_EPOCH};
|
||||
let t = UNIX_EPOCH + Duration::new(1_700_000_000, 123_456_789);
|
||||
let t = UNIX_EPOCH + Duration::new(1_700_000_000, 123_456_700);
|
||||
assert_eq!(
|
||||
MTime::from_system_time(t),
|
||||
MTime {
|
||||
secs: 1_700_000_000,
|
||||
nanos: 123_456_789
|
||||
nanos: 123_456_700
|
||||
}
|
||||
);
|
||||
assert_eq!(
|
||||
|
||||
+641
-137
@@ -53,7 +53,7 @@ use crate::daemon::control::{
|
||||
WireErrorKind, feature,
|
||||
};
|
||||
use crate::daemon::duplex::{Duplex, Halves};
|
||||
use crate::host::{Host, SharedHost, WatchSub};
|
||||
use crate::host::{Host, SearchHit, SharedHost, WatchSub};
|
||||
|
||||
/// Ceiling on threads one connection's pool will grow to.
|
||||
///
|
||||
@@ -144,6 +144,21 @@ impl Services {
|
||||
#[derive(Default)]
|
||||
pub struct AttachRegistry {
|
||||
live: Mutex<Vec<Live>>,
|
||||
/// Held across *both* tables for the length of one handover.
|
||||
///
|
||||
/// A takeover moves two things that live in different places: this
|
||||
/// registry's handles, and the `WorkspaceStore`'s record. Each is
|
||||
/// internally locked, and that is not enough — two clients attaching to one
|
||||
/// workspace at the same moment can each win a different table, after which
|
||||
/// the store names a session the registry has already evicted and no
|
||||
/// `detach` can ever clear it, because the token no longer matches. From
|
||||
/// then on the workspace reports a takeover against a client that
|
||||
/// disconnected hours ago.
|
||||
///
|
||||
/// Coarse on purpose: attach and detach happen once per workspace opened or
|
||||
/// closed, so serializing them costs nothing worth measuring. Always the
|
||||
/// outermost lock of the two, and never held while writing to a peer.
|
||||
handover: Mutex<()>,
|
||||
}
|
||||
|
||||
struct Live {
|
||||
@@ -182,6 +197,11 @@ struct Evicted {
|
||||
}
|
||||
|
||||
impl AttachRegistry {
|
||||
/// Take the handover lock. See [`AttachRegistry::handover`].
|
||||
fn handover(&self) -> std::sync::MutexGuard<'_, ()> {
|
||||
self.handover.lock().unwrap_or_else(|e| e.into_inner())
|
||||
}
|
||||
|
||||
/// Who holds `workspace` — `(token, hostname)`. Diagnostics and tests.
|
||||
pub fn holder(&self, workspace: &str) -> Option<(String, String)> {
|
||||
self.locked()
|
||||
@@ -244,6 +264,14 @@ impl AttachRegistry {
|
||||
evicted
|
||||
}
|
||||
|
||||
/// Forget `workspace` whoever holds it — the workspace itself is gone.
|
||||
///
|
||||
/// Unconditional, unlike [`AttachRegistry::release`]: a delete is not one
|
||||
/// session giving something up, it is the thing ceasing to exist.
|
||||
fn forget_workspace(&self, workspace: &str) {
|
||||
self.locked().retain(|l| l.workspace != workspace);
|
||||
}
|
||||
|
||||
/// Release `workspace`, but only if `conn` still holds it. `false` means it
|
||||
/// had already been taken over, which is success as far as the caller is
|
||||
/// concerned — and the reason releasing is conditional at all.
|
||||
@@ -372,6 +400,7 @@ where
|
||||
sink: Arc::clone(&sink),
|
||||
inflight: Mutex::new(HashMap::new()),
|
||||
watches: Mutex::new(HashMap::new()),
|
||||
deferred_watches: Mutex::new(HashMap::new()),
|
||||
next_watch: AtomicU64::new(1),
|
||||
pool: Pool::new(),
|
||||
workspaces: services.workspaces.clone(),
|
||||
@@ -520,13 +549,21 @@ fn attach_workspace(
|
||||
dedicated: bool,
|
||||
) -> io::Result<Option<String>> {
|
||||
let store = conn.workspaces()?;
|
||||
let displaced = store.attach(
|
||||
workspace,
|
||||
Attachment::new(conn.holder.token.clone(), conn.holder.hostname.clone()),
|
||||
);
|
||||
let evicted = conn
|
||||
.attachments
|
||||
.claim(workspace, conn.id, &conn.holder, dedicated);
|
||||
let (displaced, evicted) = {
|
||||
// Both tables move under one lock. Held only across the two moves —
|
||||
// the notice below goes out with nothing held, because writing to a
|
||||
// peer that has stopped reading must not hold up the next client's
|
||||
// attach.
|
||||
let _handover = conn.attachments.handover();
|
||||
let displaced = store.attach(
|
||||
workspace,
|
||||
Attachment::new(conn.holder.token.clone(), conn.holder.hostname.clone()),
|
||||
);
|
||||
let evicted = conn
|
||||
.attachments
|
||||
.claim(workspace, conn.id, &conn.holder, dedicated);
|
||||
(displaced, evicted)
|
||||
};
|
||||
|
||||
if let Some(evicted) = evicted {
|
||||
log::info!(
|
||||
@@ -562,6 +599,7 @@ fn attach_workspace(
|
||||
/// then tidied up must not evict the client that took over from it.
|
||||
fn detach_workspace(conn: &Arc<Conn>, workspace: &str) -> io::Result<bool> {
|
||||
let store = conn.workspaces()?;
|
||||
let _handover = conn.attachments.handover();
|
||||
let released = conn.attachments.release(workspace, conn.id);
|
||||
let forgotten = store.detach(workspace, &conn.holder.token);
|
||||
Ok(released || forgotten)
|
||||
@@ -637,7 +675,7 @@ fn run_job(conn: &Arc<Conn>, req_id: u64, req: ControlRequest, blob: Vec<u8>) {
|
||||
}
|
||||
|
||||
let wants_blob = req.returns_blob();
|
||||
let (reply, out_blob) = match run_request(conn, req, blob) {
|
||||
let (reply, out_blob) = match run_request(conn, req_id, req, blob) {
|
||||
Ok((ok, bytes)) => (ControlReply::Ok(ok), bytes),
|
||||
Err(e) => (ControlReply::Err(WireError::from_io(&e)), Vec::new()),
|
||||
};
|
||||
@@ -648,10 +686,36 @@ fn run_job(conn: &Arc<Conn>, req_id: u64, req: ControlRequest, blob: Vec<u8>) {
|
||||
// Dispatch
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Drop the hits whose paths this wire cannot carry.
|
||||
///
|
||||
/// `SearchHit::path` is the one `PathBuf` that crosses the control dialect —
|
||||
/// every path in `ControlRequest` is a `String` for exactly this reason — and
|
||||
/// `serde` refuses to serialize a `Path` that is not UTF-8 rather than
|
||||
/// converting it lossily. So one Latin-1 filename under the search roots makes
|
||||
/// the *whole* reply unencodable: before this, the client saw no reply at all
|
||||
/// and waited out its 20-second search deadline, on that keystroke and on every
|
||||
/// one after it.
|
||||
///
|
||||
/// Dropped rather than converted lossily, because the lossy form is a path that
|
||||
/// does not open — the client would be offering the user a hit it cannot act
|
||||
/// on. `LocalHost` keeps full fidelity; this is the wire's limit, applied at the
|
||||
/// wire.
|
||||
fn drop_unsendable_hits(hits: &mut Vec<SearchHit>) {
|
||||
let before = hits.len();
|
||||
hits.retain(|hit| hit.path.to_str().is_some());
|
||||
if hits.len() != before {
|
||||
log::debug!(
|
||||
"search dropped {} hit(s) whose paths are not UTF-8",
|
||||
before - hits.len()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// One request against the host. `Ok` carries the reply value and, for the
|
||||
/// methods that have one, the bulk bytes that ride the frame's blob.
|
||||
fn run_request(
|
||||
conn: &Arc<Conn>,
|
||||
req_id: u64,
|
||||
req: ControlRequest,
|
||||
blob: Vec<u8>,
|
||||
) -> io::Result<(ReplyOk, Vec<u8>)> {
|
||||
@@ -693,13 +757,14 @@ fn run_request(
|
||||
show_hidden,
|
||||
} => {
|
||||
let roots: Vec<PathBuf> = roots.iter().map(|r| p(r)).collect();
|
||||
let hits = h.search(
|
||||
let mut hits = h.search(
|
||||
&roots,
|
||||
&query,
|
||||
clamp_usize(limit),
|
||||
clamp_usize(max_dirs),
|
||||
show_hidden,
|
||||
)?;
|
||||
drop_unsendable_hits(&mut hits);
|
||||
(ReplyOk::Hits(hits), Vec::new())
|
||||
}
|
||||
|
||||
@@ -748,7 +813,7 @@ fn run_request(
|
||||
|
||||
// ----- watch ---------------------------------------------------------
|
||||
ControlRequest::WatchOpen { dirs } => {
|
||||
let id = conn.open_watch(&paths(&dirs))?;
|
||||
let id = conn.open_watch(req_id, &paths(&dirs))?;
|
||||
(ReplyOk::WatchId(id), Vec::new())
|
||||
}
|
||||
ControlRequest::WatchSet { id, dirs } => {
|
||||
@@ -788,7 +853,19 @@ fn run_request(
|
||||
ControlRequest::WorkspaceDelete { id } => {
|
||||
// Deleting what is not there is success — a delete that raced
|
||||
// another client's delete has got what it asked for.
|
||||
conn.workspaces()?.delete(&id, conn.workspace_origin)?;
|
||||
let store = conn.workspaces()?;
|
||||
{
|
||||
// The store drops its own attachment on delete; the registry
|
||||
// has to be told, and under the same lock, or the two disagree
|
||||
// with no race needed at all. Left behind, the stale `Live`
|
||||
// entry means the *next* client to attach a workspace with this
|
||||
// id evicts a session nobody displaced — and, that entry being
|
||||
// dedicated, closes its whole link, taking every other
|
||||
// workspace on it down too.
|
||||
let _handover = conn.attachments.handover();
|
||||
store.delete(&id, conn.workspace_origin)?;
|
||||
conn.attachments.forget_workspace(&id);
|
||||
}
|
||||
(ReplyOk::Unit, Vec::new())
|
||||
}
|
||||
|
||||
@@ -833,6 +910,17 @@ struct Conn {
|
||||
/// leaks nothing.
|
||||
inflight: Mutex<HashMap<u64, bool>>,
|
||||
watches: Mutex<HashMap<u64, WatchSub>>,
|
||||
/// Watch forwarders that have been set up but must not start until their
|
||||
/// `WatchId` reply has gone out, keyed by the request that opened them.
|
||||
///
|
||||
/// The forwarder writes to the same sink the reply does, and nothing
|
||||
/// ordered the two: a directory that changed in the instant it was first
|
||||
/// watched could push a batch that overtook the reply. The client files a
|
||||
/// watch id only once `call` returns, so it has no entry for that id yet
|
||||
/// and drops the batch under the unknown-id rule — and since the file tree
|
||||
/// relists only on a watch event, the change stays invisible until
|
||||
/// something else touches the directory.
|
||||
deferred_watches: Mutex<HashMap<u64, (u64, smol::channel::Receiver<Vec<PathBuf>>)>>,
|
||||
next_watch: AtomicU64,
|
||||
pool: Pool,
|
||||
/// The machine's workspace records, when this server serves them.
|
||||
@@ -868,6 +956,7 @@ impl Conn {
|
||||
/// earlier is already gone from the registry and is not touched — the exact
|
||||
/// case the store's token check exists for, seen from the other side.
|
||||
fn release_all_workspaces(&self) {
|
||||
let _handover = self.attachments.handover();
|
||||
let released = self.attachments.release_conn(self.id);
|
||||
let Some(store) = self.workspaces.as_ref() else {
|
||||
return;
|
||||
@@ -912,6 +1001,7 @@ impl Conn {
|
||||
.unwrap_or(false);
|
||||
if cancelled {
|
||||
log::trace!("control request {req_id} was cancelled; dropping its reply");
|
||||
self.start_deferred_watch(req_id, false);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -927,13 +1017,41 @@ impl Conn {
|
||||
} else {
|
||||
ControlServerMsg::Response { req_id, reply }
|
||||
};
|
||||
if let Err(e) = self.sink.send(&msg) {
|
||||
log::debug!("control reply {req_id} could not be written: {e}");
|
||||
}
|
||||
// Encoded before anything is written, so a reply this server cannot put
|
||||
// on the wire — a `SearchHit` whose path is not UTF-8, a `WorkspaceList`
|
||||
// grown past `MAX_FRAME` — becomes an error the client *receives*.
|
||||
// Dropping it instead leaves the client waiting out the request's whole
|
||||
// deadline (20s for a search, and again on the next keystroke) for a
|
||||
// reply that was never coming.
|
||||
let delivered = match msg.to_frame() {
|
||||
Ok((k, payload)) => match self.sink.send_frame(k, &payload) {
|
||||
Ok(()) => true,
|
||||
Err(e) => {
|
||||
log::debug!("control reply {req_id} could not be written: {e}");
|
||||
false
|
||||
}
|
||||
},
|
||||
Err(e) => {
|
||||
log::warn!("control reply {req_id} could not be encoded: {e}");
|
||||
let excuse = ControlServerMsg::Response {
|
||||
req_id,
|
||||
reply: ControlReply::Err(WireError::from_io(&e)),
|
||||
};
|
||||
if let Err(e) = self.sink.send(&excuse) {
|
||||
log::debug!("control error reply {req_id} could not be written: {e}");
|
||||
}
|
||||
false
|
||||
}
|
||||
};
|
||||
|
||||
// Only now, and only if the client actually learned the id: a batch
|
||||
// that overtook this reply would be dropped by a client that has no
|
||||
// entry for it yet. See `Conn::deferred_watches`.
|
||||
self.start_deferred_watch(req_id, delivered);
|
||||
}
|
||||
|
||||
/// Open a watch and start forwarding its batches as pushes.
|
||||
fn open_watch(&self, dirs: &[PathBuf]) -> io::Result<u64> {
|
||||
fn open_watch(&self, req_id: u64, dirs: &[PathBuf]) -> io::Result<u64> {
|
||||
let sub = self.host.watch(dirs)?;
|
||||
let id = self.next_watch.fetch_add(1, Ordering::Relaxed);
|
||||
// Clone the receiver before the subscription is filed away: the
|
||||
@@ -944,10 +1062,42 @@ impl Conn {
|
||||
.lock()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
.insert(id, sub);
|
||||
spawn_watch_forwarder(id, rx, Arc::clone(&self.sink));
|
||||
// Parked rather than started — see `deferred_watches`. `finish` starts
|
||||
// it once the reply carrying `id` is on the wire.
|
||||
self.deferred_watches
|
||||
.lock()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
.insert(req_id, (id, rx));
|
||||
Ok(id)
|
||||
}
|
||||
|
||||
/// Start the forwarder for a watch opened by `req_id`, if there was one.
|
||||
///
|
||||
/// Called from [`Conn::finish`] after the reply has gone out, and on the
|
||||
/// paths where no reply goes out at all — a cancelled request, or one whose
|
||||
/// reply could not be written — so a parked forwarder is never left holding
|
||||
/// a receiver nobody will read.
|
||||
fn start_deferred_watch(&self, req_id: u64, deliver: bool) {
|
||||
let parked = self
|
||||
.deferred_watches
|
||||
.lock()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
.remove(&req_id);
|
||||
let Some((id, rx)) = parked else {
|
||||
return;
|
||||
};
|
||||
if !deliver {
|
||||
// The client never learned this id, so every batch would be
|
||||
// dropped. Let the subscription go with the watch entry instead.
|
||||
self.watches
|
||||
.lock()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
.remove(&id);
|
||||
return;
|
||||
}
|
||||
spawn_watch_forwarder(id, rx, Arc::clone(&self.sink));
|
||||
}
|
||||
|
||||
fn set_watch_dirs(&self, id: u64, dirs: &[PathBuf]) -> io::Result<()> {
|
||||
let watches = self.watches.lock().unwrap_or_else(|e| e.into_inner());
|
||||
let sub = watches.get(&id).ok_or_else(|| {
|
||||
@@ -1065,11 +1215,21 @@ impl Sink {
|
||||
}
|
||||
|
||||
fn send(&self, msg: &ControlServerMsg) -> io::Result<()> {
|
||||
let (k, payload) = msg.to_frame()?;
|
||||
self.send_frame(k, &payload)
|
||||
}
|
||||
|
||||
/// The write half of [`Sink::send`], for callers that already encoded.
|
||||
///
|
||||
/// Split so that "this reply cannot be encoded" and "this link is gone" are
|
||||
/// distinguishable: only the second may have put bytes on the wire, and only
|
||||
/// the first leaves the connection well enough to answer on.
|
||||
fn send_frame(&self, k: u8, payload: &[u8]) -> io::Result<()> {
|
||||
let mut slot = self.out.lock().unwrap_or_else(|e| e.into_inner());
|
||||
let w = slot.as_mut().ok_or_else(|| {
|
||||
io::Error::new(io::ErrorKind::BrokenPipe, "control connection is closed")
|
||||
})?;
|
||||
msg.encode(w).and_then(|()| w.flush())
|
||||
crate::daemon::protocol::write_frame(&mut *w, k, payload).and_then(|()| w.flush())
|
||||
}
|
||||
|
||||
/// Drop the write half at teardown, so a worker that finishes afterwards
|
||||
@@ -1112,6 +1272,30 @@ struct PoolState {
|
||||
closed: bool,
|
||||
}
|
||||
|
||||
impl PoolState {
|
||||
/// Whether a job just queued needs a worker spawned for it.
|
||||
///
|
||||
/// Compares the backlog against the parked workers rather than asking
|
||||
/// whether *any* worker is parked. `idle` counts a worker from before it
|
||||
/// parks until after it has re-acquired the lock on its way out, so for the
|
||||
/// whole wake-up window a worker that has already been handed a job still
|
||||
/// looks free — and the `notify_one` a second submit sends in that window
|
||||
/// goes to a thread that has left the wait set, so it is lost.
|
||||
///
|
||||
/// Concretely: with `k` workers parked, a client pipelining `k+1` frames in
|
||||
/// one read — a `Git` plus `k` `ReadDir`s, which is what the file tree and
|
||||
/// the branch line produce together — got `k` of them running and left the
|
||||
/// last queued behind a `git status` that can take twenty seconds. That is
|
||||
/// the head-of-line blocking this pool is elastic in order to avoid.
|
||||
///
|
||||
/// Counting both sides makes the window harmless: the job destined for a
|
||||
/// parked-but-notified worker is still in `jobs`, so it cancels that
|
||||
/// worker out and the next job over sees no spare capacity.
|
||||
fn wants_another_worker(&self) -> bool {
|
||||
self.jobs.len() > self.idle && self.workers < MAX_WORKERS
|
||||
}
|
||||
}
|
||||
|
||||
impl Pool {
|
||||
fn new() -> Pool {
|
||||
Pool {
|
||||
@@ -1136,9 +1320,10 @@ impl Pool {
|
||||
}
|
||||
st.jobs.push_back(Box::new(job));
|
||||
|
||||
// Spawn only when nobody is parked to take it, so a steady stream of
|
||||
// requests is served by one warm worker rather than a thread per call.
|
||||
if st.idle == 0 && st.workers < MAX_WORKERS {
|
||||
// Spawn only when the backlog outruns the workers parked to take it, so
|
||||
// a steady stream of requests is served by one warm worker rather than a
|
||||
// thread per call.
|
||||
if st.wants_another_worker() {
|
||||
st.workers += 1;
|
||||
let inner = Arc::clone(&self.inner);
|
||||
match std::thread::Builder::new()
|
||||
@@ -1315,16 +1500,32 @@ mod sock {
|
||||
|
||||
/// Bind the control socket, replacing a stale one.
|
||||
///
|
||||
/// Permissions are the access boundary — a Unix socket has no other. The
|
||||
/// directory goes to 0700 and the socket to 0600 *before* anything can
|
||||
/// connect, so there is no window in which another user on the box could
|
||||
/// reach a server that answers `ReadFile` for arbitrary paths.
|
||||
/// Permissions are the access boundary — a Unix socket has no other, and
|
||||
/// what is behind it is `ReadFile`/`WriteFile`/`Git` on arbitrary paths as
|
||||
/// this user. So the socket has to be 0600 from the first instant its final
|
||||
/// name exists, which `bind` alone cannot give: `bind` creates the node at
|
||||
/// `0777 & ~umask`, and a `chmod` on the next line is a window. Under a
|
||||
/// `umask 002` — the default wherever user-private groups are configured —
|
||||
/// that window is group-connectable.
|
||||
///
|
||||
/// So the umask is tightened across the `bind` itself. Not a staging
|
||||
/// directory and a rename, which would also close the window: the staging
|
||||
/// path is longer than the final one, and `sun_path` is the one budget here
|
||||
/// with no room to spend — [`socket_path_in`] already falls back to a hashed
|
||||
/// name to stay under it.
|
||||
pub fn bind_control_socket(path: &Path) -> io::Result<UnixListener> {
|
||||
if let Some(parent) = path.parent() {
|
||||
let parent = path.parent().unwrap_or(Path::new("."));
|
||||
// `create_dir_all` and *only then* a chmod, on a directory that did not
|
||||
// exist a moment ago: `path` can be `$TTY7_CONTROL_SOCK` or the hashed
|
||||
// fallback, whose parent is `$XDG_RUNTIME_DIR` or `/tmp` — directories
|
||||
// this process does not own and must not re-permission. Tightening
|
||||
// `/tmp` to 0700 would lock every other user out of it, sticky bit and
|
||||
// all, with nothing linking the breakage back to tty7.
|
||||
if !parent.exists() {
|
||||
std::fs::create_dir_all(parent)?;
|
||||
// Best effort: a pre-existing `$XDG_RUNTIME_DIR` we do not own is
|
||||
// already 0700 by the platform's own rules, and failing here would
|
||||
// refuse to start over a permission we did not need to change.
|
||||
// Ours by construction, since it did not exist above. Best effort:
|
||||
// losing the race to another tty7 starting at the same instant
|
||||
// leaves the directory correct anyway.
|
||||
let _ = std::fs::set_permissions(parent, std::fs::Permissions::from_mode(0o700));
|
||||
}
|
||||
|
||||
@@ -1347,11 +1548,41 @@ mod sock {
|
||||
}
|
||||
}
|
||||
|
||||
let listener = UnixListener::bind(path)?;
|
||||
let listener = bind_private(path)?;
|
||||
// Belt and braces on two counts: it narrows the 0700 the umask below
|
||||
// yields to the 0600 a socket actually needs, and it covers a
|
||||
// filesystem that ignores the umask entirely (some FUSE mounts,
|
||||
// anything with a default ACL) — which would otherwise leave the socket
|
||||
// wide open with nothing to notice it.
|
||||
std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))?;
|
||||
Ok(listener)
|
||||
}
|
||||
|
||||
/// `bind`, with the umask tightened so the node is created owner-only
|
||||
/// rather than created at `0777 & ~umask` and fixed up afterwards.
|
||||
///
|
||||
/// `0o077`, not `0o177`: the umask is process-global for the length of the
|
||||
/// `bind`, and a *directory* another thread creates in that window would
|
||||
/// come out without its owner-execute bit and be unusable — the test suite
|
||||
/// found this the hard way. Masking only the group and other bits leaves
|
||||
/// anything created alongside it working while still giving the socket no
|
||||
/// group or world access, which is the whole property: connecting to a Unix
|
||||
/// socket needs write permission on it.
|
||||
pub(super) fn bind_private(path: &Path) -> io::Result<UnixListener> {
|
||||
// Serializes tty7's own binds against each other, so two of them cannot
|
||||
// interleave their save/restore and leave the umask tightened.
|
||||
static UMASK: Mutex<()> = Mutex::new(());
|
||||
let _held = UMASK.lock().unwrap_or_else(|e| e.into_inner());
|
||||
|
||||
// SAFETY: `umask` is always successful and has no preconditions; the
|
||||
// only hazard is the process-global effect, which the lock and the
|
||||
// immediate restore below bound.
|
||||
let previous = unsafe { libc::umask(0o077) };
|
||||
let bound = UnixListener::bind(path);
|
||||
unsafe { libc::umask(previous) };
|
||||
bound
|
||||
}
|
||||
|
||||
/// Serve control connections on `listener` until it fails, one thread per
|
||||
/// connection.
|
||||
pub fn serve_listener(listener: UnixListener, host: SharedHost) {
|
||||
@@ -1419,13 +1650,158 @@ pub use sock::{
|
||||
spawn_control_listener, spawn_control_listener_with,
|
||||
};
|
||||
|
||||
/// The pool is plain threads and channels, so unlike the rest of this file's
|
||||
/// tests — which need a Unix socket pair — these hold on every platform.
|
||||
#[cfg(test)]
|
||||
mod pool_tests {
|
||||
use super::*;
|
||||
use std::sync::atomic::AtomicBool;
|
||||
use std::time::Instant;
|
||||
|
||||
/// Workers are spawned only when nobody is free, so a serial stream of jobs
|
||||
/// costs one thread rather than one thread per job.
|
||||
#[test]
|
||||
fn the_pool_reuses_a_warm_worker() {
|
||||
let pool = Pool::new();
|
||||
|
||||
// A job signals from *inside* itself, so it has sent before its worker
|
||||
// has parked again. Submitting in that window is a genuine "nobody is
|
||||
// free" and legitimately spawns a second worker — so let the pool settle
|
||||
// first, and the assertion is about reuse rather than about timing.
|
||||
let settled = || {
|
||||
let deadline = Instant::now() + Duration::from_secs(5);
|
||||
while Instant::now() < deadline {
|
||||
let st = pool.inner.state.lock().unwrap();
|
||||
if st.idle == st.workers {
|
||||
return;
|
||||
}
|
||||
drop(st);
|
||||
std::thread::sleep(Duration::from_millis(1));
|
||||
}
|
||||
panic!("the pool never went idle");
|
||||
};
|
||||
|
||||
for _ in 0..50 {
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
assert!(pool.submit(move || {
|
||||
let _ = tx.send(());
|
||||
}));
|
||||
rx.recv_timeout(Duration::from_secs(5)).unwrap();
|
||||
settled();
|
||||
}
|
||||
let st = pool.inner.state.lock().unwrap();
|
||||
assert_eq!(
|
||||
st.workers, 1,
|
||||
"50 sequential jobs should not want 50 threads"
|
||||
);
|
||||
}
|
||||
|
||||
/// The spawn decision, at the three states that distinguish it from
|
||||
/// "is anybody parked?".
|
||||
///
|
||||
/// Driven directly because the state it has to get right — a worker that
|
||||
/// has been notified but has not yet re-acquired the lock, so it is counted
|
||||
/// in `idle` while the job meant for it is still counted in `jobs` — is a
|
||||
/// few instructions wide and a racing test passes against the broken rule
|
||||
/// far more often than it fails.
|
||||
#[test]
|
||||
fn a_worker_is_spawned_when_the_backlog_outruns_the_parked_workers() {
|
||||
let state = |jobs: usize, workers: usize, idle: usize| PoolState {
|
||||
jobs: (0..jobs).map(|_| Box::new(|| ()) as Job).collect(),
|
||||
workers,
|
||||
idle,
|
||||
closed: false,
|
||||
};
|
||||
|
||||
// One parked worker, one job: it is exactly the warm-worker case, and
|
||||
// spawning here is what would cost a thread per request.
|
||||
assert!(!state(1, 1, 1).wants_another_worker());
|
||||
|
||||
// One parked worker, two jobs. The second job arrived inside the
|
||||
// first's wake-up window, so `idle` still says 1 — and the old rule,
|
||||
// which only asked whether `idle == 0`, left this job queued behind a
|
||||
// request that can take twenty seconds.
|
||||
assert!(state(2, 1, 1).wants_another_worker());
|
||||
|
||||
// Nobody parked at all.
|
||||
assert!(state(1, 1, 0).wants_another_worker());
|
||||
|
||||
// And the ceiling still holds.
|
||||
assert!(!state(64, MAX_WORKERS, 0).wants_another_worker());
|
||||
}
|
||||
|
||||
/// And it does grow when work genuinely overlaps.
|
||||
#[test]
|
||||
fn the_pool_grows_for_concurrent_work() {
|
||||
let pool = Pool::new();
|
||||
let (done_tx, done_rx) = std::sync::mpsc::channel();
|
||||
let gate = Arc::new((Mutex::new(false), Condvar::new()));
|
||||
for _ in 0..8 {
|
||||
let gate = Arc::clone(&gate);
|
||||
let done_tx = done_tx.clone();
|
||||
assert!(pool.submit(move || {
|
||||
let (lock, cv) = &*gate;
|
||||
let mut open = lock.lock().unwrap();
|
||||
let _ = done_tx.send(());
|
||||
while !*open {
|
||||
open = cv.wait(open).unwrap();
|
||||
}
|
||||
}));
|
||||
}
|
||||
for _ in 0..8 {
|
||||
done_rx.recv_timeout(Duration::from_secs(5)).unwrap();
|
||||
}
|
||||
assert_eq!(pool.inner.state.lock().unwrap().workers, 8);
|
||||
|
||||
let (lock, cv) = &*gate;
|
||||
*lock.lock().unwrap() = true;
|
||||
cv.notify_all();
|
||||
pool.close();
|
||||
}
|
||||
|
||||
/// Closing drops the backlog. That is what breaks the `Conn` → `Pool` → job
|
||||
/// → `Arc<Conn>` cycle; leaving the jobs queued would keep every watch on a
|
||||
/// dead connection alive for the life of the process.
|
||||
#[test]
|
||||
fn closing_the_pool_drops_queued_work() {
|
||||
let pool = Pool::new();
|
||||
let gate = Arc::new((Mutex::new(false), Condvar::new()));
|
||||
let ran = Arc::new(AtomicBool::new(false));
|
||||
|
||||
// One job to occupy the single worker...
|
||||
let blocker = Arc::clone(&gate);
|
||||
assert!(pool.submit(move || {
|
||||
let (lock, cv) = &*blocker;
|
||||
let mut open = lock.lock().unwrap();
|
||||
while !*open {
|
||||
open = cv.wait(open).unwrap();
|
||||
}
|
||||
}));
|
||||
std::thread::sleep(Duration::from_millis(100));
|
||||
// ...and one that must never run.
|
||||
let ran2 = Arc::clone(&ran);
|
||||
pool.submit(move || ran2.store(true, Ordering::SeqCst));
|
||||
|
||||
pool.close();
|
||||
let (lock, cv) = &*gate;
|
||||
*lock.lock().unwrap() = true;
|
||||
cv.notify_all();
|
||||
std::thread::sleep(Duration::from_millis(200));
|
||||
assert!(
|
||||
!ran.load(Ordering::SeqCst),
|
||||
"a job queued at close still ran"
|
||||
);
|
||||
assert!(!pool.submit(|| {}), "a closed pool must refuse work");
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(all(test, unix))]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::daemon::control::{ControlHello, MTime, feature};
|
||||
use crate::host::local::LocalHost;
|
||||
use crate::host::remote::RemoteHost;
|
||||
use crate::host::{Entry, HostId, Meta, Output, SearchHit};
|
||||
use crate::host::{Entry, HostId, Meta, Output};
|
||||
use std::os::unix::net::UnixStream;
|
||||
use std::sync::atomic::AtomicBool;
|
||||
use std::time::Instant;
|
||||
@@ -2072,6 +2448,63 @@ mod tests {
|
||||
assert!(seen, "no watch event crossed the connection");
|
||||
}
|
||||
|
||||
/// The `WatchId` reply is on the wire before any batch for that id.
|
||||
///
|
||||
/// The ordering is structural — the forwarder is parked in
|
||||
/// `Conn::deferred_watches` and started by `finish`, so there is no
|
||||
/// interleaving left to hit. This is the end-to-end statement of that, run
|
||||
/// under churn; it does **not** reproduce the old race, whose window was a
|
||||
/// few instructions between `open_watch` returning and `finish` writing on
|
||||
/// the same thread.
|
||||
///
|
||||
/// What the race cost, when it landed: the client files a watch id only
|
||||
/// once its `call` returns, so a batch that overtook the reply hit a client
|
||||
/// with no entry for that id and was dropped under the unknown-id rule. The
|
||||
/// file tree relists only on a watch event, so that first change stayed
|
||||
/// invisible until something else touched the directory.
|
||||
#[test]
|
||||
fn a_watch_id_reaches_the_client_before_any_batch_for_it() {
|
||||
let (mut client, _ok) = raw();
|
||||
let tmp = tempfile::TempDir::new().unwrap();
|
||||
|
||||
ControlClientMsg::Request {
|
||||
req_id: 1,
|
||||
req: ControlRequest::WatchOpen {
|
||||
dirs: vec![tmp.path().to_string_lossy().into_owned()],
|
||||
},
|
||||
}
|
||||
.encode(&mut client)
|
||||
.unwrap();
|
||||
client.flush().unwrap();
|
||||
|
||||
// Churn from the moment the request is sent, so the window between the
|
||||
// watcher going live and the reply going out is a busy one.
|
||||
let churn = tmp.path().to_path_buf();
|
||||
let stop = Arc::new(AtomicBool::new(false));
|
||||
let churning = Arc::clone(&stop);
|
||||
let churner = std::thread::spawn(move || {
|
||||
let mut n = 0u32;
|
||||
while !churning.load(Ordering::SeqCst) {
|
||||
let _ = std::fs::write(churn.join(format!("f{n}")), b"x");
|
||||
n = n.wrapping_add(1);
|
||||
std::thread::sleep(Duration::from_millis(1));
|
||||
}
|
||||
});
|
||||
|
||||
// The very first frame back has to be the reply, not an event.
|
||||
let first = ControlServerMsg::read(&mut client).unwrap();
|
||||
stop.store(true, Ordering::SeqCst);
|
||||
churner.join().unwrap();
|
||||
|
||||
match first {
|
||||
ControlServerMsg::Response {
|
||||
req_id: 1,
|
||||
reply: ControlReply::Ok(ReplyOk::WatchId(_)),
|
||||
} => {}
|
||||
other => panic!("a batch overtook the WatchId reply: {other:?}"),
|
||||
}
|
||||
}
|
||||
|
||||
/// Dropping the subscription releases the *server's* watcher, not just the
|
||||
/// client's bookkeeping — otherwise every expanded directory in a long
|
||||
/// session leaks an OS watch on the remote machine.
|
||||
@@ -2191,112 +2624,6 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
// -----------------------------------------------------------------------
|
||||
// The pool
|
||||
// -----------------------------------------------------------------------
|
||||
|
||||
/// Workers are spawned only when nobody is free, so a serial stream of jobs
|
||||
/// costs one thread rather than one thread per job.
|
||||
#[test]
|
||||
fn the_pool_reuses_a_warm_worker() {
|
||||
let pool = Pool::new();
|
||||
|
||||
// A job signals from *inside* itself, so it has sent before its worker
|
||||
// has parked again. Submitting in that window is a genuine "nobody is
|
||||
// free" and legitimately spawns a second worker — so let the pool settle
|
||||
// first, and the assertion is about reuse rather than about timing.
|
||||
let settled = || {
|
||||
let deadline = Instant::now() + Duration::from_secs(5);
|
||||
while Instant::now() < deadline {
|
||||
let st = pool.inner.state.lock().unwrap();
|
||||
if st.idle == st.workers {
|
||||
return;
|
||||
}
|
||||
drop(st);
|
||||
std::thread::sleep(Duration::from_millis(1));
|
||||
}
|
||||
panic!("the pool never went idle");
|
||||
};
|
||||
|
||||
for _ in 0..50 {
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
assert!(pool.submit(move || {
|
||||
let _ = tx.send(());
|
||||
}));
|
||||
rx.recv_timeout(Duration::from_secs(5)).unwrap();
|
||||
settled();
|
||||
}
|
||||
let st = pool.inner.state.lock().unwrap();
|
||||
assert_eq!(
|
||||
st.workers, 1,
|
||||
"50 sequential jobs should not want 50 threads"
|
||||
);
|
||||
}
|
||||
|
||||
/// And it does grow when work genuinely overlaps.
|
||||
#[test]
|
||||
fn the_pool_grows_for_concurrent_work() {
|
||||
let pool = Pool::new();
|
||||
let (done_tx, done_rx) = std::sync::mpsc::channel();
|
||||
let gate = Arc::new((Mutex::new(false), Condvar::new()));
|
||||
for _ in 0..8 {
|
||||
let gate = Arc::clone(&gate);
|
||||
let done_tx = done_tx.clone();
|
||||
assert!(pool.submit(move || {
|
||||
let (lock, cv) = &*gate;
|
||||
let mut open = lock.lock().unwrap();
|
||||
let _ = done_tx.send(());
|
||||
while !*open {
|
||||
open = cv.wait(open).unwrap();
|
||||
}
|
||||
}));
|
||||
}
|
||||
for _ in 0..8 {
|
||||
done_rx.recv_timeout(Duration::from_secs(5)).unwrap();
|
||||
}
|
||||
assert_eq!(pool.inner.state.lock().unwrap().workers, 8);
|
||||
|
||||
let (lock, cv) = &*gate;
|
||||
*lock.lock().unwrap() = true;
|
||||
cv.notify_all();
|
||||
pool.close();
|
||||
}
|
||||
|
||||
/// Closing drops the backlog. That is what breaks the `Conn` → `Pool` → job
|
||||
/// → `Arc<Conn>` cycle; leaving the jobs queued would keep every watch on a
|
||||
/// dead connection alive for the life of the process.
|
||||
#[test]
|
||||
fn closing_the_pool_drops_queued_work() {
|
||||
let pool = Pool::new();
|
||||
let gate = Arc::new((Mutex::new(false), Condvar::new()));
|
||||
let ran = Arc::new(AtomicBool::new(false));
|
||||
|
||||
// One job to occupy the single worker...
|
||||
let blocker = Arc::clone(&gate);
|
||||
assert!(pool.submit(move || {
|
||||
let (lock, cv) = &*blocker;
|
||||
let mut open = lock.lock().unwrap();
|
||||
while !*open {
|
||||
open = cv.wait(open).unwrap();
|
||||
}
|
||||
}));
|
||||
std::thread::sleep(Duration::from_millis(100));
|
||||
// ...and one that must never run.
|
||||
let ran2 = Arc::clone(&ran);
|
||||
pool.submit(move || ran2.store(true, Ordering::SeqCst));
|
||||
|
||||
pool.close();
|
||||
let (lock, cv) = &*gate;
|
||||
*lock.lock().unwrap() = true;
|
||||
cv.notify_all();
|
||||
std::thread::sleep(Duration::from_millis(200));
|
||||
assert!(
|
||||
!ran.load(Ordering::SeqCst),
|
||||
"a job queued at close still ran"
|
||||
);
|
||||
assert!(!pool.submit(|| {}), "a closed pool must refuse work");
|
||||
}
|
||||
|
||||
// -----------------------------------------------------------------------
|
||||
// The socket
|
||||
// -----------------------------------------------------------------------
|
||||
@@ -2320,6 +2647,91 @@ mod tests {
|
||||
drop(listener);
|
||||
}
|
||||
|
||||
/// A directory that was already there is left alone. `$TTY7_CONTROL_SOCK`
|
||||
/// and the hashed fallback both put the socket straight into
|
||||
/// `$XDG_RUNTIME_DIR` or `/tmp`; tightening one of those to 0700 would lock
|
||||
/// every other user out of it — sticky bit and all, if this is running as
|
||||
/// root — with nothing linking the breakage back to tty7.
|
||||
#[test]
|
||||
fn binding_does_not_re_permission_a_directory_it_did_not_create() {
|
||||
use std::os::unix::fs::PermissionsExt as _;
|
||||
let dir = tempfile::TempDir::new().unwrap();
|
||||
let shared = dir.path().join("shared");
|
||||
std::fs::create_dir(&shared).unwrap();
|
||||
std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o1777)).unwrap();
|
||||
|
||||
let listener = bind_control_socket(&shared.join("s.sock")).unwrap();
|
||||
assert_eq!(
|
||||
std::fs::metadata(&shared).unwrap().permissions().mode() & 0o7777,
|
||||
0o1777,
|
||||
"a directory tty7 did not create must keep its own permissions"
|
||||
);
|
||||
drop(listener);
|
||||
}
|
||||
|
||||
/// Owner-only comes from the umask the bind runs under, not from a `chmod`
|
||||
/// after the fact. A socket that spends even one syscall at `0777 & ~umask`
|
||||
/// is one another user on the box can connect to, and what is behind it is
|
||||
/// `ReadFile` on arbitrary paths as this user.
|
||||
#[test]
|
||||
fn a_bind_is_owner_only_under_a_permissive_umask() {
|
||||
use std::os::unix::fs::PermissionsExt as _;
|
||||
let dir = tempfile::TempDir::new().unwrap();
|
||||
let path = dir.path().join("permissive.sock");
|
||||
|
||||
// SAFETY: `umask` has no preconditions. It is process-global, so an
|
||||
// unrelated test creating a file in this window sees `0` rather than
|
||||
// the developer's umask — more permissive, which nothing asserts on.
|
||||
let previous = unsafe { libc::umask(0) };
|
||||
let bound = sock::bind_private(&path);
|
||||
unsafe { libc::umask(previous) };
|
||||
|
||||
let listener = bound.unwrap();
|
||||
assert_eq!(
|
||||
std::fs::metadata(&path).unwrap().permissions().mode() & 0o077,
|
||||
0,
|
||||
"bind left a window in which another user could connect"
|
||||
);
|
||||
drop(listener);
|
||||
}
|
||||
|
||||
/// One filename this wire cannot carry must not cost the whole search.
|
||||
///
|
||||
/// Not driven through a real file: APFS rejects a non-UTF-8 name outright
|
||||
/// (`EILSEQ`), so the case cannot be staged on the machine most of this is
|
||||
/// developed on. The filter is the whole behaviour, so the filter is what
|
||||
/// is pinned — `to_frame_refuses_what_cannot_be_sent_without_writing_it`
|
||||
/// covers the other half, that such a hit really would fail to encode.
|
||||
#[test]
|
||||
fn a_filename_that_is_not_utf8_costs_only_its_own_hit() {
|
||||
use std::ffi::OsStr;
|
||||
use std::os::unix::ffi::OsStrExt as _;
|
||||
|
||||
let hit = |name: &str, path: &Path| SearchHit {
|
||||
name: name.to_string(),
|
||||
path: path.to_path_buf(),
|
||||
is_dir: false,
|
||||
ignored: false,
|
||||
};
|
||||
let mut hits = vec![
|
||||
hit("plain-needle.rs", Path::new("/home/me/plain-needle.rs")),
|
||||
// `café-needle.rs` in Latin-1: a lone 0xe9 is not valid UTF-8.
|
||||
hit(
|
||||
"caf\u{fffd}-needle.rs",
|
||||
Path::new(OsStr::from_bytes(b"/home/me/caf\xe9-needle.rs")),
|
||||
),
|
||||
hit("other-needle.rs", Path::new("/home/me/other-needle.rs")),
|
||||
];
|
||||
|
||||
drop_unsendable_hits(&mut hits);
|
||||
let names: Vec<&str> = hits.iter().map(|h| h.name.as_str()).collect();
|
||||
assert_eq!(
|
||||
names,
|
||||
["plain-needle.rs", "other-needle.rs"],
|
||||
"the representable hits must survive their neighbour"
|
||||
);
|
||||
}
|
||||
|
||||
/// A whole conformance-shaped exchange over a real listener, proving the
|
||||
/// socket path end to end: bind, connect, handshake, RPC.
|
||||
#[test]
|
||||
@@ -2823,6 +3235,98 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// Deleting a workspace clears it from *both* tables.
|
||||
///
|
||||
/// The store drops its own attachment on delete. If the registry keeps its
|
||||
/// handle, the two disagree with no race needed, and the next client to
|
||||
/// attach that id evicts a session nobody displaced — closing its whole
|
||||
/// link, since a dedicated entry takes every other workspace on that
|
||||
/// connection down with it.
|
||||
#[test]
|
||||
fn deleting_a_workspace_clears_both_attachment_tables() {
|
||||
let (services, _dir) = workspace_services();
|
||||
let registry = Arc::clone(&services.attachments);
|
||||
let store = services.workspaces.clone().unwrap();
|
||||
|
||||
let ((mut laptop, _), _l) =
|
||||
raw_hello(services.clone(), hello_for("w", "tok-laptop", "laptop"));
|
||||
await_holder(®istry, "w", "laptop");
|
||||
assert!(store.attachment("w").is_some());
|
||||
|
||||
ask(
|
||||
&mut laptop,
|
||||
1,
|
||||
ControlRequest::WorkspacePut {
|
||||
id: "w".to_string(),
|
||||
json: ws_record("w", "the workspace"),
|
||||
},
|
||||
);
|
||||
let reply = ask(
|
||||
&mut laptop,
|
||||
2,
|
||||
ControlRequest::WorkspaceDelete {
|
||||
id: "w".to_string(),
|
||||
},
|
||||
);
|
||||
assert!(
|
||||
matches!(reply, ControlReply::Ok(ReplyOk::Unit)),
|
||||
"{reply:?}"
|
||||
);
|
||||
|
||||
assert!(
|
||||
store.attachment("w").is_none(),
|
||||
"the store still names a holder for a workspace that is gone"
|
||||
);
|
||||
assert!(
|
||||
registry.holder("w").is_none(),
|
||||
"the registry still holds a workspace that is gone"
|
||||
);
|
||||
}
|
||||
|
||||
/// The store's record and the registry's handle move under **one** lock.
|
||||
///
|
||||
/// They are separate tables with separate locks, and taking them one after
|
||||
/// the other is not enough: two clients attaching the same workspace at the
|
||||
/// same instant can each win a different one, after which the store names a
|
||||
/// session the registry has already evicted. No `detach` can clear it — its
|
||||
/// token no longer matches — so from then on the workspace reports a
|
||||
/// takeover against a client that disconnected hours ago.
|
||||
///
|
||||
/// Held from the test rather than raced, because the window is a few
|
||||
/// instructions wide and a racing test passes against the broken ordering
|
||||
/// far more often than it fails. Holding the handover proves the stronger
|
||||
/// thing anyway: with it held, an attach reaches *neither* table.
|
||||
#[test]
|
||||
fn an_attach_moves_both_tables_under_one_lock() {
|
||||
let (services, _dir) = workspace_services();
|
||||
let registry = Arc::clone(&services.attachments);
|
||||
let store = services.workspaces.clone().unwrap();
|
||||
|
||||
let held = registry.handover();
|
||||
// The handshake replies before the attach, so this returns rather than
|
||||
// blocking on the lock we are holding.
|
||||
let ((_laptop, _ok), _served) =
|
||||
raw_hello(services.clone(), hello_for("w", "tok-laptop", "laptop"));
|
||||
|
||||
std::thread::sleep(Duration::from_millis(150));
|
||||
assert!(
|
||||
registry.holder("w").is_none(),
|
||||
"the registry was moved while a handover was in flight"
|
||||
);
|
||||
assert!(
|
||||
store.attachment("w").is_none(),
|
||||
"the store was moved while a handover was in flight"
|
||||
);
|
||||
|
||||
drop(held);
|
||||
await_holder(®istry, "w", "laptop");
|
||||
assert_eq!(
|
||||
store.attachment("w").map(|a| a.token).as_deref(),
|
||||
Some("tok-laptop"),
|
||||
"both tables have to name the same session once the handover is done"
|
||||
);
|
||||
}
|
||||
|
||||
/// The other half of that rule. A client holds **one connection per
|
||||
/// machine**, so closing the link on a takeover would drop windows nobody
|
||||
/// preempted; the push still goes out, the link stays up.
|
||||
|
||||
@@ -163,6 +163,7 @@ fn run_stdio(args: &[String]) -> io::Result<()> {
|
||||
{
|
||||
use std::os::unix::net::UnixStream;
|
||||
use tty7_core::daemon::duplex::StdioDuplex;
|
||||
use tty7_core::daemon::spawn;
|
||||
use tty7_core::host::local::LocalHost;
|
||||
use tty7_core::host::server;
|
||||
|
||||
@@ -201,10 +202,43 @@ fn run_stdio(args: &[String]) -> io::Result<()> {
|
||||
Err(e) if force_bridge => return Err(e),
|
||||
Err(e) => {
|
||||
log_stderr(format_args!(
|
||||
"no control server at {} ({e}); serving in this process",
|
||||
"no control server at {} ({e})",
|
||||
sock.display()
|
||||
));
|
||||
None
|
||||
// One control server per machine, started if nobody has —
|
||||
// the same rule `bridge_panes` follows one dialect over,
|
||||
// and for the same reason. Two `--stdio` sessions both
|
||||
// falling through to serving in-process would each hold
|
||||
// their own `WorkspaceStore` over the one file, and
|
||||
// `persist` writes the whole document: the second to save
|
||||
// silently drops the first's changes. Their attachment
|
||||
// registries would be separate too, which makes design
|
||||
// §10's takeover a no-op between them — both clients would
|
||||
// hold the same workspace and neither would be told.
|
||||
//
|
||||
// Not attempted when the caller named a socket: starting a
|
||||
// daemon binds the machine's default endpoint, not theirs,
|
||||
// so it would be a daemon nobody asked for and nobody uses.
|
||||
if may_start_daemon(args) {
|
||||
match spawn::ensure_running()
|
||||
.map_err(io::Error::other)
|
||||
.and_then(|()| UnixStream::connect(&sock))
|
||||
{
|
||||
Ok(s) => {
|
||||
log_stderr(format_args!("started one; bridging to it"));
|
||||
Some(s)
|
||||
}
|
||||
Err(e) => {
|
||||
log_stderr(format_args!(
|
||||
"could not start one ({e}); serving in this process"
|
||||
));
|
||||
None
|
||||
}
|
||||
}
|
||||
} else {
|
||||
log_stderr(format_args!("serving in this process"));
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
@@ -358,6 +392,17 @@ fn control_services() -> tty7_core::host::server::Services {
|
||||
}
|
||||
|
||||
/// `--flag <value>` or `--flag=<value>`, first occurrence wins.
|
||||
/// Whether a failed control probe may start the machine's daemon.
|
||||
///
|
||||
/// Only when the caller did not name a socket. `--control-sock` says "this
|
||||
/// endpoint", and `spawn::ensure_running` binds the machine's default one — so
|
||||
/// starting a daemon there would leave a process nobody asked for and nobody
|
||||
/// reaches. It is also what keeps the test suite, and any `--config-dir`
|
||||
/// isolation built on it, from spraying daemons across a developer's machine.
|
||||
fn may_start_daemon(args: &[String]) -> bool {
|
||||
flag_value(args, "--control-sock").is_none()
|
||||
}
|
||||
|
||||
fn flag_value(args: &[String], flag: &str) -> Option<String> {
|
||||
let with_eq = format!("{flag}=");
|
||||
let mut it = args.iter();
|
||||
@@ -385,3 +430,23 @@ fn apply_config_dir_arg(args: &[String]) {
|
||||
tty7_core::core::config::set_config_dir(path.into());
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn argv(args: &[&str]) -> Vec<String> {
|
||||
args.iter().map(|a| a.to_string()).collect()
|
||||
}
|
||||
|
||||
/// A named socket suppresses the daemon start, in both spellings of the
|
||||
/// flag. Without this guard every `--stdio` in the test suite that points
|
||||
/// at a temp socket would start a real daemon on the developer's machine.
|
||||
#[test]
|
||||
fn a_named_control_socket_suppresses_starting_a_daemon() {
|
||||
assert!(may_start_daemon(&argv(&[])));
|
||||
assert!(may_start_daemon(&argv(&["--serve"])));
|
||||
assert!(!may_start_daemon(&argv(&["--control-sock", "/tmp/x.sock"])));
|
||||
assert!(!may_start_daemon(&argv(&["--control-sock=/tmp/x.sock"])));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -6,22 +6,40 @@
|
||||
//! a connection to a control server that is already running, which is the path
|
||||
//! an `ssh host tty7-server --stdio` takes on a machine with a live daemon and
|
||||
//! which no amount of `Host` conformance would exercise.
|
||||
//!
|
||||
//! Everything `--stdio` is Unix-only — the flag is refused on Windows, where a
|
||||
//! machine is reached over its own transport rather than by shipping a server
|
||||
//! onto it (contract §8). The plain argument handling below is not, and runs
|
||||
//! everywhere.
|
||||
|
||||
use std::process::{Command, Stdio};
|
||||
|
||||
#[cfg(unix)]
|
||||
use std::io;
|
||||
#[cfg(unix)]
|
||||
use std::path::PathBuf;
|
||||
use std::process::{Child, Command, Stdio};
|
||||
#[cfg(unix)]
|
||||
use std::process::Child;
|
||||
#[cfg(unix)]
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
#[cfg(unix)]
|
||||
use tty7_core::daemon::control::{ControlHello, LinkShutdown};
|
||||
#[cfg(unix)]
|
||||
use tty7_core::host::Host;
|
||||
#[cfg(unix)]
|
||||
use tty7_core::host::local::LocalHost;
|
||||
#[cfg(unix)]
|
||||
use tty7_core::host::remote::RemoteHost;
|
||||
#[cfg(unix)]
|
||||
use tty7_core::host::server;
|
||||
|
||||
const EXE: &str = env!("CARGO_BIN_EXE_tty7-server");
|
||||
|
||||
#[cfg(unix)]
|
||||
struct ServerProcess(Mutex<Option<Child>>);
|
||||
|
||||
#[cfg(unix)]
|
||||
impl LinkShutdown for ServerProcess {
|
||||
fn shutdown_link(&self) -> io::Result<()> {
|
||||
if let Some(mut c) = self.0.lock().unwrap_or_else(|e| e.into_inner()).take() {
|
||||
@@ -33,6 +51,7 @@ impl LinkShutdown for ServerProcess {
|
||||
}
|
||||
|
||||
/// Start `tty7-server --stdio <args>` and connect a `RemoteHost` to its pipes.
|
||||
#[cfg(unix)]
|
||||
fn stdio_child(args: &[&str]) -> io::Result<Arc<RemoteHost>> {
|
||||
let mut child = Command::new(EXE)
|
||||
.arg("--stdio")
|
||||
@@ -49,6 +68,7 @@ fn stdio_child(args: &[&str]) -> io::Result<Arc<RemoteHost>> {
|
||||
}
|
||||
|
||||
/// A control server on a temp socket, for the bridge to reach.
|
||||
#[cfg(unix)]
|
||||
fn listening_server(dir: &tempfile::TempDir) -> PathBuf {
|
||||
let sock = dir.path().join("control.sock");
|
||||
let listener = server::bind_control_socket(&sock).unwrap();
|
||||
@@ -63,6 +83,7 @@ fn listening_server(dir: &tempfile::TempDir) -> PathBuf {
|
||||
/// stream carries belongs to the client and the server at the far end, and a
|
||||
/// bridge with an opinion about the protocol would become a third party to a
|
||||
/// negotiation it is not qualified to join.
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn the_bridge_carries_a_whole_session() {
|
||||
let dir = tempfile::TempDir::new().unwrap();
|
||||
@@ -92,6 +113,7 @@ fn the_bridge_carries_a_whole_session() {
|
||||
/// `--bridge` with nowhere to bridge to fails rather than quietly serving
|
||||
/// itself. An operator who asked for the bridge is telling us a server exists;
|
||||
/// silently becoming that server would fork the machine's state in two.
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn an_explicit_bridge_with_no_server_fails() {
|
||||
let dir = tempfile::TempDir::new().unwrap();
|
||||
@@ -105,6 +127,7 @@ fn an_explicit_bridge_with_no_server_fails() {
|
||||
|
||||
/// With neither flag, `--stdio` probes: nothing listening means serve here, so a
|
||||
/// machine that has never run a daemon is still reachable over ssh.
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn the_default_mode_serves_when_nothing_is_listening() {
|
||||
let dir = tempfile::TempDir::new().unwrap();
|
||||
@@ -117,6 +140,7 @@ fn the_default_mode_serves_when_nothing_is_listening() {
|
||||
|
||||
/// ...and something listening means bridge to it, so a second `--stdio` session
|
||||
/// joins the machine's existing server instead of standing up a rival.
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn the_default_mode_bridges_when_a_server_is_listening() {
|
||||
let dir = tempfile::TempDir::new().unwrap();
|
||||
@@ -130,6 +154,7 @@ fn the_default_mode_bridges_when_a_server_is_listening() {
|
||||
}
|
||||
|
||||
/// Contradictory flags are refused rather than one silently winning.
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn serve_and_bridge_together_are_refused() {
|
||||
let out = Command::new(EXE)
|
||||
|
||||
@@ -22,6 +22,11 @@
|
||||
//! isolation: no case can be explained by another's leftover state, a hung
|
||||
//! server fails exactly one case, and a crash names the behaviour that caused it.
|
||||
|
||||
// Unix-only: every case here is a `--stdio` child, and `--stdio` is refused on
|
||||
// Windows by design — a Windows machine is reached over its own transport, not
|
||||
// by shipping a server onto it (contract §8).
|
||||
#![cfg(unix)]
|
||||
|
||||
use std::io;
|
||||
use std::path::Path;
|
||||
use std::process::{Child, Command, Stdio};
|
||||
|
||||
@@ -19,6 +19,10 @@
|
||||
//! another's leftovers and nothing here can touch the developer's real
|
||||
//! `~/.local/share/tty7/workspaces.json`.
|
||||
|
||||
// Unix-only, for the same reason as `stdio_conformance.rs`: the server under
|
||||
// test is a `--stdio` child, and two of the cases stand up a control socket.
|
||||
#![cfg(unix)]
|
||||
|
||||
use std::io;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::process::{Child, Command, Stdio};
|
||||
|
||||
@@ -1802,6 +1802,13 @@ impl RemoteTerminal {
|
||||
|
||||
/// Send one workspace-scoped request and return the daemon's reply.
|
||||
///
|
||||
/// How long a workspace-addressed request waits for the daemon.
|
||||
///
|
||||
/// Generous, because behind it is an SSH round trip to the workspace's own
|
||||
/// machine and possibly a connection being established — but finite, which
|
||||
/// is the point.
|
||||
const WORKSPACE_OP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
|
||||
|
||||
/// The counterpart of the `pane_id`-addressed helpers above for a pane that
|
||||
/// lives on a *remote workspace*: there is no pane on the local daemon to
|
||||
/// name, so the request carries the workspace and a secret-free spec naming
|
||||
@@ -1813,6 +1820,17 @@ impl RemoteTerminal {
|
||||
/// empty list.
|
||||
pub fn on_workspace(req: WorkspaceRequest) -> anyhow::Result<DaemonMsg> {
|
||||
let mut stream = connect()?;
|
||||
// Bounded, because the daemon's answer is not just its own work: it
|
||||
// resolves the workspace's SSH connection and, for the forward ops,
|
||||
// waits for the *server* to acknowledge a `cancel_tcpip_forward`. On a
|
||||
// box that has gone unreachable — lid closed, VPN dropped, which is
|
||||
// exactly when someone reaches for Stop Workspace — that acknowledgement
|
||||
// never comes. Without a deadline this read parks forever, and the
|
||||
// thread with it.
|
||||
//
|
||||
// Best effort: a transport that will not take a timeout degrades to the
|
||||
// old unbounded read rather than failing the request outright.
|
||||
let _ = stream.set_read_timeout(Some(Self::WORKSPACE_OP_TIMEOUT));
|
||||
ClientMsg::OnWorkspace(Box::new(req)).encode(&mut stream)?;
|
||||
match DaemonMsg::read(&mut stream)? {
|
||||
DaemonMsg::Error(msg) => Err(anyhow::anyhow!(msg)),
|
||||
|
||||
+26
-5
@@ -1263,7 +1263,17 @@ impl Tty7App {
|
||||
// An empty workspace has nothing to come back to, so it is dropped
|
||||
// outright instead of accumulating as a blank row in the picker —
|
||||
// every `New Workspace` the user closes without using would leave one.
|
||||
if self.tabs.is_empty() {
|
||||
//
|
||||
// Unless the emptiness is *this client's* ignorance rather than the
|
||||
// machine's answer. `claimable_session` deliberately opens a remote
|
||||
// workspace empty when its machine cannot be reached, so a window
|
||||
// opened while the box was asleep and then closed — there was nothing
|
||||
// in it to work on — would take the entry with it: its `RemoteRef`, its
|
||||
// cached layout and its geometry, while its panes are still running
|
||||
// over there. Nothing would reconnect it and nothing would offer it
|
||||
// again; the only way back is re-adding the machine by hand.
|
||||
let answered = WorkspaceStore::machine_is_connected(cx, self.workspace);
|
||||
if self.tabs.is_empty() && answered {
|
||||
WorkspaceStore::remove(cx, self.workspace);
|
||||
} else {
|
||||
WorkspaceStore::close_window(cx, self.workspace);
|
||||
@@ -1297,10 +1307,21 @@ impl Tty7App {
|
||||
else {
|
||||
return;
|
||||
};
|
||||
let left = route.teardown();
|
||||
if !left.is_empty() {
|
||||
log::warn!("{} forwards survived a workspace teardown", left.len());
|
||||
}
|
||||
// Off the UI thread. `teardown` dials the daemon, which resolves the
|
||||
// workspace's SSH connection and waits for the server to acknowledge a
|
||||
// `cancel_tcpip_forward` — on a machine that has gone unreachable, which
|
||||
// is exactly when someone reaches for Stop Workspace, that never comes
|
||||
// back inside the request timeout. `ForwardRoute::list` is already
|
||||
// backgrounded for the same reason; this was the one that was not, and
|
||||
// it ran while the window was being torn down.
|
||||
cx.background_executor()
|
||||
.spawn(async move {
|
||||
let left = route.teardown();
|
||||
if !left.is_empty() {
|
||||
log::warn!("{} forwards survived a workspace teardown", left.len());
|
||||
}
|
||||
})
|
||||
.detach();
|
||||
}
|
||||
|
||||
/// Stop a workspace — kill its sessions and close its window — confirming
|
||||
|
||||
@@ -200,6 +200,12 @@ fn agent_icon(path: &str) -> Option<&'static [u8]> {
|
||||
// silhouette, so this is lobehub/lobe-icons' square transcription (MIT),
|
||||
// drawn for exactly this avatar use. Its notice rides in the SVG.
|
||||
"icons/agents/grok.svg" => include_bytes!("../../assets/icons/agents/grok.svg"),
|
||||
// Pi's own mark, from pi.dev — the one file that arrives with theme
|
||||
// logic attached (`logo-auto.svg` carries a `prefers-color-scheme`
|
||||
// style block). The geometry is kept as published and rescaled to the
|
||||
// 24x24 grid; the CSS is dropped, since these avatars are tinted by the
|
||||
// app and no other mark here brings a stylesheet. Details in the SVG.
|
||||
"icons/agents/pi.svg" => include_bytes!("../../assets/icons/agents/pi.svg"),
|
||||
_ => return None,
|
||||
};
|
||||
Some(bytes)
|
||||
|
||||
+32
-1
@@ -41,7 +41,7 @@ use gpui_component::{
|
||||
};
|
||||
|
||||
use crate::ui::app::Tty7App;
|
||||
use crate::ui::host_ops::{HostOps, MTime, WatchSub};
|
||||
use crate::ui::host_ops::{HostOps, MTime, SharedHost, WatchSub};
|
||||
|
||||
/// Refuse to open files larger than this: the component's code editor is rated
|
||||
/// to ~50K lines, and a multi-megabyte blob is almost never what a terminal
|
||||
@@ -165,6 +165,15 @@ pub(crate) struct EditorPanelState {
|
||||
/// and a server-side watcher recreated every time a file is opened or
|
||||
/// closed. `Arc` because `set_dirs` is itself a host call.
|
||||
watch: Option<Arc<WatchSub>>,
|
||||
/// The host `watch` was opened against, kept so a subscription is never
|
||||
/// reused across a different one.
|
||||
///
|
||||
/// A `HostId` is not enough to tell them apart: reconnecting removes the
|
||||
/// dead `RemoteHost` and inserts a fresh one under the *same* id, so the id
|
||||
/// matches while the `ControlClient` behind the old subscription is gone.
|
||||
/// Compared by pointer, which distinguishes both that and an outright
|
||||
/// switch to another machine.
|
||||
watch_host: Option<SharedHost>,
|
||||
/// A subscription is being opened; keeps a burst of opens from asking for
|
||||
/// one each.
|
||||
watch_opening: bool,
|
||||
@@ -217,6 +226,7 @@ impl EditorPanelState {
|
||||
.detach();
|
||||
Self {
|
||||
watch: None,
|
||||
watch_host: None,
|
||||
watch_opening: false,
|
||||
watch_busy: false,
|
||||
watch_dirty: false,
|
||||
@@ -434,6 +444,25 @@ impl Tty7App {
|
||||
return;
|
||||
};
|
||||
|
||||
// Same rule as the file tree's: a subscription belongs to the host that
|
||||
// opened it. A reconnect inserts a fresh `RemoteHost` under the same
|
||||
// `HostId`, so the id matches while the `ControlClient` behind this
|
||||
// subscription is gone — `set_dirs` then fails, is warned and dropped,
|
||||
// and nothing opens a new one. The cost here is quieter and worse than
|
||||
// a stale tree: external-change detection is what stops a save
|
||||
// clobbering an edit made on the other side.
|
||||
if !self
|
||||
.editor
|
||||
.watch_host
|
||||
.as_ref()
|
||||
.is_some_and(|opened_with| Arc::ptr_eq(opened_with, &host))
|
||||
{
|
||||
self.editor.watch = None;
|
||||
self.editor.watch_host = None;
|
||||
self.editor.watch_busy = false;
|
||||
self.editor.watch_dirty = false;
|
||||
}
|
||||
|
||||
if let Some(sub) = self.editor.watch.clone() {
|
||||
if self.editor.watch_busy {
|
||||
self.editor.watch_dirty = true;
|
||||
@@ -462,6 +491,7 @@ impl Tty7App {
|
||||
return;
|
||||
}
|
||||
self.editor.watch_opening = true;
|
||||
let opened_host = Arc::clone(&host);
|
||||
let opened_with = self.editor.watched_dirs.clone();
|
||||
HostOps::run(
|
||||
host,
|
||||
@@ -481,6 +511,7 @@ impl Tty7App {
|
||||
};
|
||||
let events = sub.events().clone();
|
||||
app.editor.watch = Some(sub);
|
||||
app.editor.watch_host = Some(opened_host);
|
||||
cx.spawn(async move |app, cx| {
|
||||
while let Ok(batch) = events.recv().await {
|
||||
let ok = app.update(cx, |app, _cx| {
|
||||
|
||||
@@ -211,6 +211,15 @@ pub(crate) struct FileTreeState {
|
||||
/// triangle. `Arc` because `set_dirs` is itself a host call and has to be
|
||||
/// handed to the background executor.
|
||||
watch: Option<Arc<WatchSub>>,
|
||||
/// The host `watch` was opened against, kept so a subscription is never
|
||||
/// reused across a different one.
|
||||
///
|
||||
/// A `HostId` is not enough to tell them apart: reconnecting removes the
|
||||
/// dead `RemoteHost` and inserts a fresh one under the *same* id, so the id
|
||||
/// matches while the `ControlClient` behind the old subscription is gone.
|
||||
/// Compared by pointer, which distinguishes both that and an outright
|
||||
/// switch to another machine.
|
||||
watch_host: Option<SharedHost>,
|
||||
/// A subscription is being opened. Without this, render would ask for one
|
||||
/// per frame until the first answer lands.
|
||||
watch_opening: bool,
|
||||
@@ -255,6 +264,7 @@ impl FileTreeState {
|
||||
})
|
||||
.detach();
|
||||
Self {
|
||||
watch_host: None,
|
||||
children: ByHost::default(),
|
||||
loads: InFlight::default(),
|
||||
stale: HashSet::new(),
|
||||
@@ -283,6 +293,25 @@ impl FileTreeState {
|
||||
fn sync_watch(&mut self, host: SharedHost, dirs: HashSet<PathBuf>, cx: &mut Context<Tty7App>) {
|
||||
self.watched = dirs;
|
||||
let want: Vec<PathBuf> = self.watched.iter().cloned().collect();
|
||||
// A subscription belongs to the host that opened it. Reconnecting drops
|
||||
// the dead `RemoteHost` and inserts a fresh one under the same
|
||||
// `HostId`, and adopting another workspace can change the host outright
|
||||
// — in both cases the subscription here is over a `ControlClient` that
|
||||
// is gone. `set_dirs` on it then fails with `ConnectionReset`, which is
|
||||
// warned and dropped, and nothing ever opens a new one: after the first
|
||||
// reconnect of a remote workspace the tree stops seeing changes made on
|
||||
// the far side for the rest of the window's life.
|
||||
if !self
|
||||
.watch_host
|
||||
.as_ref()
|
||||
.is_some_and(|opened_with| Arc::ptr_eq(opened_with, &host))
|
||||
{
|
||||
// Dropping the subscription is what unsubscribes, on both sides.
|
||||
self.watch = None;
|
||||
self.watch_host = None;
|
||||
self.watch_busy = false;
|
||||
self.watch_dirty = false;
|
||||
}
|
||||
if let Some(sub) = self.watch.clone() {
|
||||
if self.watch_busy {
|
||||
self.watch_dirty = true;
|
||||
@@ -316,6 +345,7 @@ impl FileTreeState {
|
||||
}
|
||||
self.watch_opening = true;
|
||||
let host_id = host.id();
|
||||
let opened_host = Arc::clone(&host);
|
||||
let opened_with = self.watched.clone();
|
||||
HostOps::run(
|
||||
host,
|
||||
@@ -337,6 +367,7 @@ impl FileTreeState {
|
||||
// independent of the subscription the state holds.
|
||||
let events = sub.events().clone();
|
||||
app.file_tree.watch = Some(sub);
|
||||
app.file_tree.watch_host = Some(opened_host);
|
||||
cx.spawn(async move |app, cx| {
|
||||
while let Ok(batch) = events.recv().await {
|
||||
let ok = app.update(cx, |app, _cx| {
|
||||
|
||||
+178
-5
@@ -45,7 +45,7 @@ use std::borrow::Borrow;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::hash::Hash;
|
||||
|
||||
use gpui::{App, AppContext as _, Context, Window};
|
||||
use gpui::{App, Context, Window};
|
||||
use gpui_component::WindowExt as _;
|
||||
|
||||
// The host vocabulary, re-exported so a view imports everything it needs from
|
||||
@@ -56,6 +56,156 @@ pub use tty7_core::host::{
|
||||
Entry, Host, HostId, MTime, Meta, Output, SearchHit, SharedHost, WatchSub,
|
||||
};
|
||||
|
||||
/// Where blocking [`Host`] calls actually run.
|
||||
///
|
||||
/// **Not gpui's background executor.** On Linux that is a fixed pool of
|
||||
/// `available_parallelism().max(2)` worker threads with no separate blocking
|
||||
/// tier, so N stalled host calls on an N-core client occupy every worker there
|
||||
/// is. Everything else that uses `background_executor` then queues behind them
|
||||
/// — including the reconnect in `remote_workspace::launch_attempt`, which is
|
||||
/// the one thing that would clear the stall. Expanding a subtree on a link that
|
||||
/// has gone silent is enough: one call per directory, each parked for its
|
||||
/// deadline (5s for a `ReadDir`, 30s for a `ReadFile`). macOS is far less
|
||||
/// exposed, since libdispatch grows its global queues when their threads block,
|
||||
/// which is why this does not show up in development.
|
||||
///
|
||||
/// Elastic and its own: a thread per concurrent call, reused while warm and
|
||||
/// retired after [`LINGER`], capped at [`MAX_THREADS`]. A `Host` call is
|
||||
/// user-driven — a directory expanded, a file opened — not per-frame, so the
|
||||
/// steady state is one or two threads.
|
||||
mod blocking {
|
||||
use std::collections::VecDeque;
|
||||
use std::sync::{Arc, Condvar, Mutex, OnceLock};
|
||||
use std::time::Duration;
|
||||
|
||||
type Job = Box<dyn FnOnce() + Send + 'static>;
|
||||
|
||||
/// Ceiling on threads. Deliberately well above any core count: these are
|
||||
/// parked on a socket rather than competing for CPU, and what has to fit is
|
||||
/// the number of host calls in flight — one per expanded directory in a
|
||||
/// burst, plus whatever the editor and the git probes are doing.
|
||||
const MAX_THREADS: usize = 64;
|
||||
|
||||
/// How long an idle worker waits for more work before retiring.
|
||||
const LINGER: Duration = Duration::from_secs(30);
|
||||
|
||||
struct Inner {
|
||||
state: Mutex<State>,
|
||||
wake: Condvar,
|
||||
}
|
||||
|
||||
struct State {
|
||||
jobs: VecDeque<Job>,
|
||||
threads: usize,
|
||||
/// Workers parked in `wait_timeout`, counted from before they park
|
||||
/// until after they have re-acquired the lock on the way out.
|
||||
idle: usize,
|
||||
}
|
||||
|
||||
impl State {
|
||||
/// Whether a job just queued needs a thread spawned for it.
|
||||
///
|
||||
/// Compares the backlog against the parked workers rather than asking
|
||||
/// whether *any* worker is parked: `idle` still counts a worker that
|
||||
/// has been handed a job but has not yet woken, and the job meant for
|
||||
/// it is still in `jobs`, so counting both sides cancels the window out.
|
||||
fn wants_another_thread(&self) -> bool {
|
||||
self.jobs.len() > self.idle && self.threads < MAX_THREADS
|
||||
}
|
||||
}
|
||||
|
||||
fn pool() -> &'static Arc<Inner> {
|
||||
static POOL: OnceLock<Arc<Inner>> = OnceLock::new();
|
||||
POOL.get_or_init(|| {
|
||||
Arc::new(Inner {
|
||||
state: Mutex::new(State {
|
||||
jobs: VecDeque::new(),
|
||||
threads: 0,
|
||||
idle: 0,
|
||||
}),
|
||||
wake: Condvar::new(),
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
/// Queue `job`. Never refuses: a dropped job is a `Host` call whose caller
|
||||
/// waits forever, and at this cap the backlog is a better failure than that.
|
||||
pub(super) fn submit(job: impl FnOnce() + Send + 'static) {
|
||||
let inner = pool();
|
||||
let mut st = inner.state.lock().unwrap_or_else(|e| e.into_inner());
|
||||
st.jobs.push_back(Box::new(job));
|
||||
if st.wants_another_thread() {
|
||||
st.threads += 1;
|
||||
let spawned = Arc::clone(inner);
|
||||
match std::thread::Builder::new()
|
||||
.name("tty7-host-op".into())
|
||||
.spawn(move || worker(spawned))
|
||||
{
|
||||
Ok(_) => return,
|
||||
Err(e) => {
|
||||
// Out of threads: leave the job for whoever is already
|
||||
// running. With nobody at all there is no one to run it, so
|
||||
// run it here — blocking this caller, which is the lesser
|
||||
// harm against never answering.
|
||||
st.threads -= 1;
|
||||
log::warn!("could not start a host-op thread: {e}");
|
||||
if st.threads == 0
|
||||
&& let Some(job) = st.jobs.pop_back()
|
||||
{
|
||||
drop(st);
|
||||
job();
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
drop(st);
|
||||
inner.wake.notify_one();
|
||||
}
|
||||
|
||||
fn worker(inner: Arc<Inner>) {
|
||||
loop {
|
||||
let job = {
|
||||
let mut st = inner.state.lock().unwrap_or_else(|e| e.into_inner());
|
||||
loop {
|
||||
if let Some(job) = st.jobs.pop_front() {
|
||||
break job;
|
||||
}
|
||||
st.idle += 1;
|
||||
let (guard, timeout) = inner
|
||||
.wake
|
||||
.wait_timeout(st, LINGER)
|
||||
.unwrap_or_else(|e| e.into_inner());
|
||||
st = guard;
|
||||
st.idle -= 1;
|
||||
if timeout.timed_out() && st.jobs.is_empty() {
|
||||
st.threads -= 1;
|
||||
return;
|
||||
}
|
||||
}
|
||||
};
|
||||
job();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Run `f` on the blocking pool and await its result.
|
||||
///
|
||||
/// `None` means the job was dropped without running, which happens only when
|
||||
/// the process is going down — the caller lands nothing rather than inventing
|
||||
/// an answer.
|
||||
async fn off_thread<T, F>(f: F) -> Option<T>
|
||||
where
|
||||
T: Send + 'static,
|
||||
F: FnOnce() -> T + Send + 'static,
|
||||
{
|
||||
let (tx, rx) = smol::channel::bounded(1);
|
||||
blocking::submit(move || {
|
||||
let _ = tx.send_blocking(f());
|
||||
});
|
||||
rx.recv().await.ok()
|
||||
}
|
||||
|
||||
/// The GPUI-facing facade over [`Host`].
|
||||
///
|
||||
/// A unit struct rather than a value: there is no per-instance state, and
|
||||
@@ -84,7 +234,9 @@ impl HostOps {
|
||||
// optimizer sees through after the first.
|
||||
tty7_core::host::register_ui_thread();
|
||||
cx.spawn(async move |this, cx| {
|
||||
let out = cx.background_spawn(async move { f(&*host) }).await;
|
||||
let Some(out) = off_thread(move || f(&*host)).await else {
|
||||
return;
|
||||
};
|
||||
let _ = this.update(cx, |view, cx| land(view, out, cx));
|
||||
})
|
||||
.detach();
|
||||
@@ -114,7 +266,9 @@ impl HostOps {
|
||||
{
|
||||
tty7_core::host::register_ui_thread();
|
||||
cx.spawn(async move |_this, cx| {
|
||||
let out = cx.background_spawn(async move { f(&*host) }).await;
|
||||
let Some(out) = off_thread(move || f(&*host)).await else {
|
||||
return;
|
||||
};
|
||||
cx.update(|cx| land(cx, out));
|
||||
})
|
||||
.detach();
|
||||
@@ -131,7 +285,9 @@ impl HostOps {
|
||||
{
|
||||
tty7_core::host::register_ui_thread();
|
||||
cx.spawn_in(window, async move |this, cx| {
|
||||
let out = cx.background_spawn(async move { f(&*host) }).await;
|
||||
let Some(out) = off_thread(move || f(&*host)).await else {
|
||||
return;
|
||||
};
|
||||
let _ = this.update_in(cx, |view, window, cx| land(view, out, window, cx));
|
||||
})
|
||||
.detach();
|
||||
@@ -492,7 +648,24 @@ mod gpui_tests {
|
||||
// The pane goes away before either answer can land — a tab closed
|
||||
// while its probe was in flight.
|
||||
drop(pane);
|
||||
cx.background_executor.run_until_parked();
|
||||
|
||||
// Pumped rather than parked once: a `Host` call runs on `HostOps`' own
|
||||
// thread pool, not on gpui's executor, so `run_until_parked` has
|
||||
// nothing to wait for until the answer has already crossed back. The
|
||||
// deadline is generous because what is being timed is a closure that
|
||||
// returns a constant.
|
||||
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
|
||||
while detached.load(Ordering::SeqCst) == 0 && std::time::Instant::now() < deadline {
|
||||
cx.background_executor.run_until_parked();
|
||||
std::thread::sleep(std::time::Duration::from_millis(1));
|
||||
}
|
||||
// Both were submitted together, so the detached one landing means the
|
||||
// view-scoped one has had its chance too. A few more turns, in case it
|
||||
// is merely slower.
|
||||
for _ in 0..20 {
|
||||
cx.background_executor.run_until_parked();
|
||||
std::thread::sleep(std::time::Duration::from_millis(1));
|
||||
}
|
||||
|
||||
assert_eq!(
|
||||
detached.load(Ordering::SeqCst),
|
||||
|
||||
+25
-3
@@ -339,11 +339,17 @@ mod tests {
|
||||
assert_ne!(normal.data, attention.data);
|
||||
}
|
||||
|
||||
/// The avatar renders for a branded agent, an unbranded (bot-fallback)
|
||||
/// agent, and with/without the status dot.
|
||||
/// Every avatar renders, with and without the status dot. Run over the
|
||||
/// whole roster rather than one branded and one fallback agent, because
|
||||
/// this is the only test that puts the bundled SVGs through resvg: the
|
||||
/// asset-source test next door proves the bytes resolve, not that they
|
||||
/// parse into visible geometry. Vendor marks arrive in whatever shape the
|
||||
/// vendor publishes — stylesheets, nested groups, features usvg quietly
|
||||
/// drops — and a mark that parses to nothing shows up as a bare accent
|
||||
/// disc, which nothing else here would catch.
|
||||
#[test]
|
||||
fn agent_avatar_renders_brand_and_fallback() {
|
||||
for agent in [CLIAgent::Claude, CLIAgent::Qwen] {
|
||||
for agent in CLIAgent::ALL {
|
||||
let idle = agent_avatar(agent, AgentStatus::Idle).unwrap();
|
||||
let waiting = agent_avatar(agent, AgentStatus::Waiting).unwrap();
|
||||
assert_eq!((idle.width(), idle.height()), (32, 32));
|
||||
@@ -351,6 +357,22 @@ mod tests {
|
||||
assert_eq!(idle.pixel(0, 0).unwrap().alpha(), 0);
|
||||
// …and the center is covered (disc + glyph).
|
||||
assert!(idle.pixel(16, 16).unwrap().alpha() > 0);
|
||||
// The glyph actually drew something. The disc under it is a flat
|
||||
// accent fill, so every opaque pixel shares one colour unless the
|
||||
// white mark landed on top — one colour means resvg handed back an
|
||||
// empty canvas, which is what a silently-unsupported SVG feature
|
||||
// looks like from here.
|
||||
let shades: std::collections::HashSet<_> = idle
|
||||
.pixels()
|
||||
.iter()
|
||||
.filter(|p| p.alpha() == 0xFF)
|
||||
.map(|p| (p.red(), p.green(), p.blue()))
|
||||
.collect();
|
||||
assert!(
|
||||
shades.len() > 1,
|
||||
"{} rendered as a bare disc — its glyph drew nothing",
|
||||
agent.display_name()
|
||||
);
|
||||
// The status dot changes the bottom-right corner.
|
||||
assert_ne!(idle.data(), waiting.data());
|
||||
}
|
||||
|
||||
+73
-28
@@ -280,20 +280,41 @@ pub fn menu_order(cx: &App) -> Vec<(WorkspaceId, bool)> {
|
||||
/// [`pane_liveness`](crate::terminal::pane_liveness): the prompt states an exact
|
||||
/// number about an irreversible action, so it wants a fresh count, not one that
|
||||
/// may be ten seconds old. This runs on a click, not on a frame.
|
||||
pub fn live_pane_count(cx: &App, workspace: WorkspaceId) -> Option<usize> {
|
||||
/// What [`live_pane_count`] needs from the app, gathered on the UI thread so the
|
||||
/// count itself does not have to run there.
|
||||
pub struct PaneCountQuery {
|
||||
route: crate::terminal::PaneRoute,
|
||||
claimed: Vec<u64>,
|
||||
}
|
||||
|
||||
/// Read the inputs for [`live_pane_count`]. Cheap; UI thread only.
|
||||
pub fn pane_count_query(cx: &App, workspace: WorkspaceId) -> Option<PaneCountQuery> {
|
||||
let ws = WorkspaceStore::all(cx).get(workspace)?;
|
||||
let claimed = ws.pane_ids();
|
||||
Some(PaneCountQuery {
|
||||
// Routed to the workspace's own machine: a remote workspace's pane ids
|
||||
// mean nothing to this computer's daemon, so asking it would count
|
||||
// whichever *local* panes happen to hold those numbers and put a "3
|
||||
// running sessions will be ended" warning on a workspace that has none.
|
||||
route: crate::ui::remote_workspace::pane_route_for(cx, workspace),
|
||||
claimed: ws.pane_ids(),
|
||||
})
|
||||
}
|
||||
|
||||
/// **Blocking. Never call this on the UI thread.**
|
||||
///
|
||||
/// For a remote route this dials the workspace's machine — an SSH handshake if
|
||||
/// nothing is pooled — and a WSL one can go as far as installing the server
|
||||
/// binary. `guard_off_ui` makes a UI-thread call a debug-build abort rather
|
||||
/// than a dropped frame, which is what it did when this was reached straight
|
||||
/// from the Stop/Delete action handler.
|
||||
pub fn live_pane_count(q: &PaneCountQuery) -> Option<usize> {
|
||||
let PaneCountQuery { route, claimed } = q;
|
||||
if claimed.is_empty() {
|
||||
return Some(0);
|
||||
}
|
||||
// One short-lived connection, only when there is something to ask about —
|
||||
// the picker renders far more often than a workspace is closed. Routed to
|
||||
// the workspace's own machine: a remote workspace's pane ids mean nothing
|
||||
// to this computer's daemon, so asking it would count whichever *local*
|
||||
// panes happen to hold those numbers and put a "3 running sessions will be
|
||||
// ended" warning on a workspace that has none.
|
||||
let route = crate::ui::remote_workspace::pane_route_for(cx, workspace);
|
||||
match crate::terminal::RemoteTerminal::try_list_panes_on(&route) {
|
||||
// the picker renders far more often than a workspace is closed.
|
||||
match crate::terminal::RemoteTerminal::try_list_panes_on(route) {
|
||||
Ok(panes) => {
|
||||
let alive: std::collections::HashSet<u64> = panes
|
||||
.into_iter()
|
||||
@@ -369,33 +390,57 @@ fn confirm_destructive(
|
||||
verb: &'static str,
|
||||
act: fn(&mut App, WorkspaceId),
|
||||
) {
|
||||
let live = live_pane_count(cx, workspace);
|
||||
let name = WorkspaceStore::all(cx)
|
||||
.get(workspace)
|
||||
.map(|w| w.display_name())
|
||||
.unwrap_or_else(|| "this workspace".to_string());
|
||||
// Only a machine that *answered* zero licenses skipping the prompt. An
|
||||
// unreachable one is the case most likely to still have work in it.
|
||||
if live == Some(0) && verb == "Stop" {
|
||||
act(cx, workspace);
|
||||
return;
|
||||
}
|
||||
let detail = destructive_detail(live, verb);
|
||||
// Title Case, like every other prompt title in the app — this one used to
|
||||
// lowercase "workspace" while its siblings read "Close Window?" /
|
||||
// "Quit and Stop Daemon?".
|
||||
let answer = window.prompt(
|
||||
gpui::PromptLevel::Warning,
|
||||
&format!("{verb} Workspace \u{201c}{name}\u{201d}?"),
|
||||
Some(&detail),
|
||||
&["Cancel", verb],
|
||||
cx,
|
||||
);
|
||||
let query = pane_count_query(cx, workspace);
|
||||
let handle = window.window_handle();
|
||||
|
||||
cx.spawn(async move |cx| {
|
||||
// The count dials the workspace's machine, so it does not belong on the
|
||||
// UI thread — on a remote route that is an SSH handshake, and on a WSL
|
||||
// one it can go as far as installing the server. Reached straight from
|
||||
// the action handler, it was a `guard_off_ui` abort in a debug build
|
||||
// and a window frozen for the length of a connect in a release one.
|
||||
let live = match query {
|
||||
Some(q) => {
|
||||
cx.background_spawn(async move { live_pane_count(&q) })
|
||||
.await
|
||||
}
|
||||
None => None,
|
||||
};
|
||||
|
||||
// Only a machine that *answered* zero licenses skipping the prompt. An
|
||||
// unreachable one is the case most likely to still have work in it.
|
||||
if live == Some(0) && verb == "Stop" {
|
||||
let _ = cx.update(|cx| act(cx, workspace));
|
||||
return;
|
||||
}
|
||||
|
||||
let detail = destructive_detail(live, verb);
|
||||
// Title Case, like every other prompt title in the app — this one used
|
||||
// to lowercase "workspace" while its siblings read "Close Window?" /
|
||||
// "Quit and Stop Daemon?".
|
||||
let Ok(answer) = handle.update(cx, |_, window, cx| {
|
||||
window.prompt(
|
||||
gpui::PromptLevel::Warning,
|
||||
&format!("{verb} Workspace \u{201c}{name}\u{201d}?"),
|
||||
Some(&detail),
|
||||
&["Cancel", verb],
|
||||
cx,
|
||||
)
|
||||
}) else {
|
||||
// The window went away while we were asking its machine. Nothing to
|
||||
// confirm against, and acting unprompted is exactly what this path
|
||||
// exists to prevent.
|
||||
return;
|
||||
};
|
||||
|
||||
// Index 1 == the verb button; Cancel and a dismissed prompt both leave
|
||||
// the workspace alone.
|
||||
if let Ok(1) = answer.await {
|
||||
cx.update(|cx| act(cx, workspace));
|
||||
let _ = cx.update(|cx| act(cx, workspace));
|
||||
}
|
||||
})
|
||||
.detach();
|
||||
|
||||
Reference in New Issue
Block a user