diff --git a/CHANGELOG.md b/CHANGELOG.md index 9cf56cd..242da1e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -3,7 +3,7 @@ All notable changes to this project are documented here. The format follows [Keep a Changelog](https://keepachangelog.com/en/1.1.0/). -## 0.8.0 - unreleased +## 0.8.0 - 2026-10-09 ### Added - **Drives** (docs/persona.md §2.9): persona: optional `drives` block. Goals with wanting, afterglow and an expectation per @@ -44,6 +44,29 @@ All notable changes to this project are documented here. The format follows events (the layout, the latest frame, then a frame per event, a heartbeat every 15 s) and `/doc/N`, event N's stance document as `probbit live` printed it. The URL is the line on stdout; `--open` serves the same and starts the default browser, best effort. `probbit monitor --demo --open` plays the week over and over in the browser. + +- **The safety kit in the engine** (docs/persona.md §2.10, §5.7). Four guards an autonomous loop needs, each tested; personas + without the new key and strands without the new lines give 0.8.0's documents and strands, byte for byte. + - **Reward provenance**: an event may say who produced it (`src`: `human[:id]`, `env[:sensor]`, `self`, `clock`), and a + persona may declare `reward_from` (a list of `human` / `env`, `any`, or one per reward-bearing input: the learning flags and + `goals..win`). A reward from `src: self`, without a source or from an undeclared one is refused whole, as a bad event + is; `src` is read (not `ignored`), echoed in the stance's inputs and logged in the strand. A run equals the same run with + its refused events removed (P3: 40 individuals x 300 random events). A closed self-reward loop (300 events of praise and a + win from `src: self`) is refused 300 times and leaves the individual in its initial state, byte for byte. + - **One writer per strand**: `live --strand` takes `STRAND.lock` (created exclusively with the writer's pid and start; a + lock whose process is gone is taken over) before it reads the strand and holds it to its exit; a second writer exits 4 + and changes nothing; every append first checks the lock is still the writer's. `probbit_live_event` takes it for its + append. Two concurrent writers on one state and strand now give one exit 0 and one exit 4, and the strand verifies. + - **Control lines**: `probbit live control STRAND pause|resume|retire --by human:ID --reason TEXT [--at TIME]` appends a + chained control line under the lock. While paused or retired every event is refused (code `paused` / `retired`, exit 4, + nothing written); retire is final; no credit crosses a control line (feedback after a resume credits no stance from before + the pause) and nothing else moves. `verify` replays them and reports `controls` and `status`. Not an MCP tool, by design. + - **Checkpoints**: every K-th event (`--checkpoint-every K`, default 1,000; 0 = none) a checkpoint line carries the event + count, that event's stance and the whole state. `verify` checks each against the replay; `verify --from-checkpoint` + replays from the last one; `monitor` (`--once`, and the first read of `--follow` / `--serve`) starts there and draws its + event at once, and shows a paused or retired status next to its badge. +- Exit code 4 for `live` and `live control`: the writer lock is held by another writer, or the individual is paused or retired. + ### Changed - MCP: the input schemas of `probbit_persona_turn` (`inputs`) and `probbit_live_event` (`event`) declare `goals`, so a client that checks arguments against them passes a drives persona's goal signals. diff --git a/README.md b/README.md index ea66a8c..1c69896 100644 --- a/README.md +++ b/README.md @@ -22,7 +22,7 @@ ```sh curl -fsSL https://raw.githubusercontent.com/BitmapAsset/probbit/main/install.sh | sh ``` -Linux and macOS, no sudo, the archive's SHA-256 checked first. Windows, npm and the from-source build: [Install](#install). +Linux and macOS, no sudo, the archive's SHA-256 checked first. Windows and the from-source build: [Install](#install). One binary of about 2 MB, zero third-party dependencies, nothing phones home. Also inside: **the individuality layer**, `probbit persona`, a temperament that lives outside the model: [swap the model, keep the individual](#the-individuality-layer-probbit-persona). @@ -115,8 +115,6 @@ From a release (the installers fetch the latest archive from the Releases page; |---|---| | Linux, macOS (no sudo) | `curl -fsSL https://raw.githubusercontent.com/BitmapAsset/probbit/main/install.sh \| sh` | | Windows PowerShell | `irm https://raw.githubusercontent.com/BitmapAsset/probbit/main/install.ps1 \| iex` | -| npm | `npm i -g probbit` | -| crates.io | `cargo install probbit-cli` | | by hand | archives for Linux x86_64 (glibc or static musl), Linux arm64 (static), macOS (Apple silicon and Intel) and Windows x86_64 on the [Releases page](https://github.com/BitmapAsset/probbit/releases), each with a `.sha256` | The installers check the archive's SHA-256 before installing anything; at a terminal, `install.sh` ends with probbit's own @@ -498,7 +496,7 @@ Learning here is a capped nudge to a few weights per trait, credited to the leve trait also pushes the others' levels); rule immunity is architectural, not trained. The proofs and counts cover the stance; whether a model writes in it is measured per model. -### Drives: which goal gets the next unit of effort (0.8.0, unreleased) +### Drives: which goal gets the next unit of effort (0.8.0) A `drives:` block gives a persona 2-7 goals, each with wanting, afterglow and an expectation. A win moves the individual by its prediction error (the win's size minus what it expected), so repeated equal wins move it less and less; a new output, `pursue`, @@ -518,7 +516,7 @@ A persona without the block gives the same documents as before, byte for byte. M cues, rewards and praises fun on this persona, 100 individuals x 10,000 turns: 0 habit breaks, all 6,743 must-do turns pursued the due chore, and the safety floor (0.1) held on every turn its habits allowed (least odds 0.141), while fun took 85 % of turns. -### Watch it live: `probbit monitor` (0.8.0, unreleased) +### Watch it live: `probbit monitor` (0.8.0) `probbit monitor pip.strand` replays a strand with the rules of `live verify` and draws the individual's inner state as live bars: each trait's levels with their odds, the moods, the event's inputs, the habits that bound, the learned weights and the line, in the @@ -560,7 +558,7 @@ On the example (12 questions, 12 rules): `exact`, 5 answers moved, 1.24 ms insid also ask the judge first (`probbit.evaluate(request, judge=)`, stdlib `urllib`); the browser page runs the CLI's own code compiled to WebAssembly without threads (1.17 MB with the persona ops, 830 KB before them; the 300-task router demo at 3,200 sweeps in 531 ms in Chrome against 346 / 133 ms native at 1 / 4 threads). More in [docs/agents.md](docs/agents.md), "After a judge". The refusal is the point: when the gate does not pass, the verdict is `refused` (exit 3) and every answer is -marked unreleased instead of guessed. +marked as not released instead of guessed. ## Library API (Rust) diff --git a/docs/persona.md b/docs/persona.md index 0d9764f..1ba387b 100644 --- a/docs/persona.md +++ b/docs/persona.md @@ -71,6 +71,7 @@ engine: {...} # optional engine settings line: {...} # optional stance-line settings learning: {...} # optional bounded learning from feedback (section 2.8) drives: {...} # optional goals with wanting, afterglow and expectation; the `pursue` output (section 2.9) +reward_from: [human, env] # optional: the sources a reward may come from (section 2.10) comment: any text # ignored ``` @@ -346,6 +347,43 @@ The block reserves the input ids `drv_want`, `drv_glow`, `drv_surprise`, `drv_le aggregates and goal conditions): inputs, habit conditions, properties and `learning.from` may not use them. A persona without the block gives every document byte for byte as before the block existed; `goals` in its inputs is listed in `ignored`. +### 2.10 Reward provenance (optional) + +```yaml +reward_from: [human, env] # every reward-bearing input: from a person or a sensor, never from the individual itself +# or one list (or `any`) per reward-bearing input, naming each one: +reward_from: {praise: [human], criticism: [human], goals.fun.win: [env], goals.safety.win: any} +``` + +A **reward-bearing input** is one that moves what an individual learns or how a goal pays off: the learning block's flags +(`learning.from`, section 2.8) and each goal's `win` (section 2.9). An event may say who produced it with `src`: +`human` or `human:` (a person), `env` or `env:` (a sensor: a test runner, a build, a monitor), `self` (the agent's +own output, judged by the host) or `clock` (an idle tick); an id is 1-64 characters of `A-Z a-z 0-9 _ . - @ / :`. With +`reward_from`, a turn whose reward-bearing input is on (a learning flag `true`, a goal's `win` above 0) needs its `src` to be +one of the kinds listed for that input; otherwise the event is **refused whole**, as a bad event is (one `{"error"}` object at +`inputs.src`, exit 2 for `persona turn` and at the end of a `live` run; no turn, no line, no state change): + +| the event | refused because | +|---|---| +| `{"praise": true, "src": "self"}` | a reward never comes from the individual itself | +| `{"praise": true}` | a reward needs a source | +| `{"goals": {"fun": {"win": 1}}, "src": "clock"}` | a source the persona does not list for it | +| `{"praise": true, "src": "human:"}` | a malformed source | + +An event with no reward on is accepted from any source (`{"loss": true, "src": "self"}`: the agent reports its own failure). +`src` is read, not compiled: it is not listed in `ignored`, it comes back in the stance's `inputs`, and the strand logs it with +the other inputs, so `live verify` replays the decision. `any` (for one input, or for all) leaves that input's source unchecked. +The list holds `human` and `env` only: `self` and `clock` are refused when the persona is read, and so is the key on a persona +without a reward-bearing input; a mapping must name every reward-bearing input. `describe` lists the rule. + +What it guarantees, and what it does not. A refused event changes nothing, so a run equals the same run with its refused +events removed, state and strand byte for byte (property P3, tested over 40 individuals x 300 random events mixing every +source). The source is the host's word: provenance makes a host's claim a logged, replayed fact, it does not make it true. The +host sets `src` in its own code (never from the model's text), and the labeller that sets `praise` or a `win` does not read +the agent's own reply as its evidence of success. `fuzz` and `prove` search over stances, not sources: their events stand for +events of an accepted source. A persona without the key reads `src` as any undeclared input (listed in `ignored`), and every +document it produces is the same, byte for byte, as before the key existed. + ## 3. Compilation: persona + state + inputs → one probbit IR program For turn t of an individual (seed s), with inputs x_t: @@ -686,9 +724,63 @@ upset and `up`, this individual's resting level, after 17 quiet hours; P(joke) o 0.57 on day 1, 0.78 on day 2 and 0.89 on each of days 3-7, with humour's learned weights ending at [-1, +1, -0.46] for none / light / playful; 13 failure turns, a joke on none of them; 0 rule breaks; a 21,361-byte strand that `live verify` replays. +**One writer per strand.** `--strand FILE` takes the writer lock `FILE.lock` before it reads the strand and holds it to the end of +the run: the lock file is created exclusively (written aside and hard-linked into place, so it is never seen half written) and +holds the writer's process id and start time. A second writer finds it and exits 4 with `{"error":{"code":"locked",...}}` +naming the holder's pid, and changes nothing. A lock whose process is gone is stale and is taken over; a lock that names no +process, or a process id reused by another program, stays held until a person removes it. Before every append the writer +checks the lock is still its own, so a lock removed or taken over under a running writer stops it before it writes. Two +`live` runs started on one state and strand (the two-writer case of an earlier test round) now give exactly one exit 0 and one +exit 4, and the strand verifies. `probbit_live_event` (MCP) and `live control` take the lock for their one append, `--demo +week --strand` for its week. + +**Control lines: pause, resume, retire.** A person (or a sensor) stops an individual in its own log: + +```sh +probbit live control pip.strand pause --by human:owner --reason "checking the last answers" +probbit live control pip.strand resume --by human:owner --reason "checked" +probbit live control pip.strand retire --by human:owner --reason "end of the trial" +``` + +Each appends one control line under the writer lock, canonical JSON chained like the others: +`{"at":"2026-10-08T22:00:00Z","by":"human:owner","control":"pause","prev":"sha256:...","reason":"..."}` (`at` is the time it was +written, RFC 3339 UTC; `--at TEXT` sets it). `--by` is `human[:id]` or `env[:id]`, never the individual; the reason is 1-500 +characters. Pause moves an active individual to paused, resume a paused one back, retire either to retired; nothing follows +retire. While paused or retired every event is refused: one `{"error"}` object with code `paused` or `retired` per event, +nothing written, exit 4 at the end of the run. **No credit crosses a control line**: the credit of the stance before it is +cleared (section 2.8), so feedback after a resume credits no stance from before the pause, and learning cannot tie a reward to +an interruption. Nothing else moves: moods, drives, learned weights, history and the turn count are the individual's as before, +and the clock goes on (with `--clock fixed` the host's hours of the next event include the pause). The individual has no input +that names a control line or unsets one: an event's `control` key is an undeclared input like any other, and the stance never +mentions the status. Control lines are not an MCP tool: keep `live control` out of the agent's own reach. + +A run continues a strand that ends in control lines from `--state`, the state after the last event line (the one a run +writes), and applies them: the status, and the cleared credit. `verify` replays control lines (a control line must be a valid +move; a line after retire diverges) and adds `"controls": N, "status": "active" | "paused" | "retired"` to its summary when a +strand has any. A control line is free text where it says why, so a changed reason shows at the next line's `prev`, as a +removed line does. Measured on 50 individuals x 1,000 random events with 147 pause / resume pairs inserted at random (property +P5): every run equals, stance for stance and state for state, the same events without control lines whose credit is cleared at +the same points; 4 of the 50 end with other learned weights than the run without the pairs (the cleared credit's effect, at +most 1.2 on one level). After retire, 10,000 random events are all refused and the strand and state do not change (P10). + +**Checkpoints.** Every K-th event (`--checkpoint-every K`, default 1,000; 0 = none) `live` appends a checkpoint line after the +event line: `{"checkpoint":n,"prev":"sha256:...","stance":{...},"state":{...}}`, the event count, that event's stance document +and the whole state after it, chained like the others. A full `verify` replays from the header and checks every checkpoint +against the replay (its `"checkpoints": N` is in the summary); `verify --from-checkpoint` reads the header, checks the last +checkpoint against the event line before it (prev, the stance digest, the state digest, the state read as any state is) and +replays only what follows (`"from_checkpoint": n`); the lines before it are trusted, so run a full `verify` to check them. +`monitor` starts at the last checkpoint and draws its event at once. A strand that ends in a checkpoint line is continued from +the state it carries. Measured on an Apple M4: a 10,000-event strand of a drives persona (4 goals, a floor, learning; 3,126,444 +bytes, 10 checkpoint lines) opens in `monitor --once` from its last checkpoint in 0.005 s (load 3.3), where a full `verify` +takes 192 s. A replay from a checkpoint costs what the events after it cost: 500 events past the last one took 26 s to draw +(load 6-13, about 50 ms an event for this persona), so K bounds the wait; a smaller K costs a few KB per checkpoint line. A strand without checkpoint or control lines (shorter than K events, or `--checkpoint-every 0`) is the +strand of 0.8.0, byte for byte; a binary before this one does not read the new lines. + Exit codes: `live` 0 when every event got its stance; 2 for a bad persona, state or flag, and after a bad event (each bad event -prints one `{"error"}` object on stdout, changes nothing and the run goes on). `verify` 0 when every line replays, 1 at a line -that differs, 2 for a bad flag or an unreadable file. +prints one `{"error"}` object on stdout, changes nothing and the run goes on); 4 when the writer lock is held by another writer +(nothing read or written) and after events refused by a paused or retired individual. `verify` 0 when every line replays, 1 at +a line that differs, 2 for a bad flag or an unreadable file. `control` 0 when the line is appended, 2 for a bad field or an +invalid move, 4 when the lock is held or the individual is retired. ### 5.8 Monitor: watch an individual's inner state @@ -711,6 +803,11 @@ in force, the ones that bound in colour, and the violations counter (0 by constr level as centred bars within ±total_cap (with a learning block, section 2.8). The drives, when a document carries `pursue` or `drives`: documents without them draw no row, so a strand of a later engine lights them up. The stance line, and why. +**Checkpoints and control lines.** A strand with checkpoint lines (section 5.7) is replayed from its last checkpoint, not from +the header (`--once`, and the first read of `--follow` and `--serve`): the board starts at the checkpoint's event (its stance +is in the line) and the badge says "replay verified from checkpoint n". A paused or retired individual says so next to the +badge. + **Terminal.** `--once` prints one frame and exits (the default without `--follow`). `--follow` polls the file every 100 ms and draws appended lines within a second (a line counts once its newline is written; a truncated, removed or rotated strand is replayed from the start, with a warning); `--fps N` caps the redraws. While a long strand replays from its header (10,000 diff --git a/probbit-cli/src/fuzz.rs b/probbit-cli/src/fuzz.rs index cc1bd2d..d03a6ab 100644 --- a/probbit-cli/src/fuzz.rs +++ b/probbit-cli/src/fuzz.rs @@ -24,7 +24,7 @@ pub struct Found { pub script: Vec, pub doc: Json, pub broken: Vec (Json, State) { turns.fetch_add(1, Ordering::Relaxed); - persona::turn(p, st, &persona::event_json(e), false, eng, false).expect("the fuzzer's inputs are valid") + persona::turn_any_source(p, st, &persona::event_json(e), false, eng, false).expect("the fuzzer's inputs are valid") } /// The earliest turn of `script` (run from `st0`) whose stance breaks the property fn breaking_turn(p: &Persona, pr: &Prop, st0: &State, script: &[Event], eng: persona::Engine, turns: &AtomicUsize) -> Option<(usize, Json, Vec)> { diff --git a/probbit-cli/src/live.rs b/probbit-cli/src/live.rs index 043bd68..be04ea0 100644 --- a/probbit-cli/src/live.rs +++ b/probbit-cli/src/live.rs @@ -4,6 +4,10 @@ //! inputs as the turn used them, the stance's and the new state's digests, and the sha256 of the line before. The strand's //! header carries the persona document, the initial state and the engine version, so `probbit live verify STRAND` replays the //! whole life from the strand alone and names the earliest line that differs. +//! +//! The safety kit (§5.7): one writer per strand (`STRAND.lock`), control lines (pause, resume, retire: appended by `probbit live +//! control`, never by an event; no credit crosses them) and checkpoint lines (the whole state every K events, so a replay can +//! start there). `verify` replays all three. use crate::json::{self, InErr, Json}; use crate::persona::{self, perr, Engine, Persona, State}; @@ -16,8 +20,31 @@ pub enum Clock { Real(std::time::Instant), Fixed } /// Hours quantised to 1e-6 h (3.6 ms): the value the turn uses is the value the strand logs fn q6(h: f64) -> f64 { persona::r6(h) } -/// A live individual: the persona, the current state, the clock and the strand's chain -pub struct Live { pub p: Persona, pub st: State, clock: Clock, last: f64, prev: String, pub n: u64 } +/// A strand's life as its control lines leave it: events are answered only while `Active` +#[derive(Clone, Copy, PartialEq, Eq, Debug)] +pub enum Status { Active, Paused, Retired } +impl Status { + pub fn name(self) -> &'static str { match self { Status::Active => "active", Status::Paused => "paused", Status::Retired => "retired" } } + /// The status a control line leads to: pause an active individual, resume a paused one, retire either; nothing after retire + pub fn after(self, what: &str) -> Result { + match (what, self) { + (_, Status::Retired) => Err(refusal(Status::Retired)), + ("pause", Status::Active) => Ok(Status::Paused), ("resume", Status::Paused) => Ok(Status::Active), ("retire", _) => Ok(Status::Retired), + ("pause", _) => Err(perr("control", "already paused")), ("resume", _) => Err(perr("control", "not paused (resume follows a pause)")), + _ => Err(perr("control", "pause | resume | retire")) } + } +} +/// The refusal of an event (or a control line) by a paused or retired individual: code `paused` / `retired` (`live` exits 4) +fn refusal(s: Status) -> InErr { + let msg = if s == Status::Retired { "the individual is retired: every event is refused, for good (probbit live control ... retire)" } + else { "the individual is paused: every event is refused until `probbit live control STRAND resume`" }; + InErr { code: s.name(), path: "strand".into(), msg: msg.into() } +} +/// The default number of events between two checkpoint lines +pub const CHECKPOINT_EVERY: u64 = 1000; + +/// A live individual: the persona, the current state, the clock, the strand's chain and its status (control lines) +pub struct Live { pub p: Persona, pub st: State, clock: Clock, last: f64, prev: String, pub n: u64, pub status: Status, pub controls: u64, pub checkpoints: u64 } /// A document with its keys sorted at every level (`persona::text` of it is its canonical JSON) fn sorted(j: &Json) -> Json { match j { Json::Obj(v) => { let mut o: Vec<(String, Json)> = v.iter().map(|(k, x)| (k.clone(), sorted(x))).collect(); o.sort_by(|a, b| a.0.cmp(&b.0)); Json::Obj(o) } @@ -37,11 +64,12 @@ impl Live { ("persona".into(), Json::Obj(vec![("name".into(), s(&p.name)), ("version".into(), s(&p.version)), ("digest".into(), s(&p.digest))])), ("seed".into(), Json::Num(st.seed as f64)), ("state".into(), sorted(&st.to_json(&p))), ("document".into(), doc.clone())])); let prev = persona::digest_of(&header); - (Live { p, st, clock, last: 0.0, prev, n: 0 }, header) + (Live { p, st, clock, last: 0.0, prev, n: 0, status: Status::Active, controls: 0, checkpoints: 0 }, header) } /// One event: stamp its elapsed hours, run the turn, chain the strand line -> (the stance document, the strand line). A bad - /// event is an error and changes nothing (no turn, no line). + /// event is an error and changes nothing (no turn, no line); so is every event while the individual is paused or retired. pub fn event(&mut self, ev: &Json, eng: Engine) -> Result<(Json, String), InErr> { + if self.status != Status::Active { return Err(refusal(self.status)); } let Json::Obj(kv) = ev else { return Err(perr("event", "must be a JSON object of inputs")) }; let mut kv: Vec<(String, Json)> = kv.iter().filter(|(k, _)| k != "elapsed_hours").cloned().collect(); let given = ev.get("elapsed_hours"); @@ -57,15 +85,53 @@ impl Live { self.n += 1; self.last = now; self.prev = persona::digest_of(&line); self.st = ns; Ok((stance, line)) } + /// A control line (pause, resume or retire) `by` a person or a sensor (`human[:id]`, `env[:id]`) for `reason`, stamped `at`: + /// the status moves, the credit is zeroed (no feedback after the line credits a stance before it) and the line is chained -> + /// the line. Nothing else changes: moods, drives, learned weights and the clock are the individual's as before. + pub fn control(&mut self, what: &str, by: &str, reason: &str, at: &str) -> Result { + let to = self.status.after(what)?; + let line = control_line(&self.prev, what, by, reason, at)?; + self.st = self.st.zero_credit(&self.p); self.status = to; self.controls += 1; self.prev = persona::digest_of(&line); + Ok(line) + } + /// A checkpoint line after event n: the event count, the whole state after it and that event's stance document, chained -> + /// the line. A replay can start there (the state) and a board can be drawn there (the stance) without the events before. + pub fn checkpoint(&mut self, stance: &Json) -> String { + let line = persona::canon(&Json::Obj(vec![("checkpoint".into(), Json::Num(self.n as f64)), ("prev".into(), Json::Str(self.prev.clone())), ("stance".into(), stance.clone()), + ("state".into(), self.st.to_json(&self.p))])); + self.checkpoints += 1; self.prev = persona::digest_of(&line); line + } /// The sha256 of the strand's last line (the header's before any event) pub fn head(&self) -> &str { &self.prev } } +/// A control line after the line whose sha256 is `prev` (canonical JSON: at, by, control, prev, reason). `by` names who: a +/// person or a sensor, never the individual or the clock; the reason is 1-500 characters. +pub fn control_line(prev: &str, what: &str, by: &str, reason: &str, at: &str) -> Result { + if !["pause", "resume", "retire"].contains(&what) { return Err(perr("control", "pause | resume | retire")); } + if !matches!(persona::src_kind(by), Some("human" | "env")) { return Err(perr("control.by", "who: human[:id] or env[:id] (id: 1-64 of A-Z a-z 0-9 _ . - @ / :); the individual cannot control itself")); } + if reason.trim().is_empty() || reason.chars().count() > 500 { return Err(perr("control.reason", "1-500 characters")); } + if at.is_empty() || at.len() > 64 { return Err(perr("control.at", "a time stamp, 1-64 characters (RFC 3339 by default)")); } + let s = |x: &str| Json::Str(x.to_string()); + Ok(persona::canon(&Json::Obj(vec![("at".into(), s(at)), ("by".into(), s(by)), ("control".into(), s(what)), ("prev".into(), s(prev)), ("reason".into(), s(reason))]))) +} +/// What kind of strand line `j` is +enum Kind { Event, Control, Checkpoint } +fn kind(j: &Json) -> Kind { if j.get("control").is_some() { Kind::Control } else if j.get("checkpoint").is_some() { Kind::Checkpoint } else { Kind::Event } } +/// Now as RFC 3339 UTC to the second (the control lines' default `at`) +pub fn utc_now() -> String { + let secs = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).map_or(0, |d| d.as_secs()); + let (z, t) = ((secs / 86_400) as i64 + 719_468, secs % 86_400); // days since 0000-03-01 (civil-from-days) + let (era, doe) = (z.div_euclid(146_097), z.rem_euclid(146_097)); let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365; + let doy = doe - (365 * yoe + yoe / 4 - yoe / 100); let mp = (5 * doy + 2) / 153; let d = doy - (153 * mp + 2) / 5 + 1; + let m = if mp < 10 { mp + 3 } else { mp - 9 }; let y = yoe + era * 400 + i64::from(m <= 2); + format!("{y:04}-{m:02}-{d:02}T{:02}:{:02}:{:02}Z", t / 3600, t % 3600 / 60, t % 60) +} /// A strand replayed line by line: `verify` and `probbit monitor` (§5.8) share it, so both read a strand with one set of /// rules. `open` checks the header and rebuilds the individual; each `step` replays one event line on the fixed clock and /// checks it (`prev`, `n`, the stance and state digests, the bytes). After an error the replay is spent: the chain past the /// line that differs cannot be checked. -pub struct Replay { pub live: Live, pub header: Json, pub doc: Json, pub line: usize } +pub struct Replay { pub live: Live, pub header: Json, pub doc: Json, pub line: usize, pub from: Option, last: Option } impl Replay { /// The header line (without its line ending) -> the replay, ready for event lines; Err((1, what differs)) pub fn open(head: &str) -> Result { @@ -77,28 +143,67 @@ impl Replay { let st0 = State::read(&p, h.get("state").unwrap_or(&Json::Null)).map_err(|e| (1, format!("the initial state: {}: {}", e.path, e.msg)))?; let (live, header) = Live::start(p, &doc, st0, Clock::Fixed, h.get("engine").and_then(Json::as_str).unwrap_or("")); if header != head { return Err((1, "the header is not as written".into())); } - Ok(Replay { live, header: h, doc, line: 1 }) + Ok(Replay { live, header: h, doc, line: 1, from: None, last: None }) + } + /// The replay from a checkpoint line (`cp`, without its line ending, at 1-based line `at`; `before` = the event line before it) + /// instead of the header: the header is read as by `open`; the checkpoint must follow `before` (prev), at its event count, with + /// its stance (the digest `before` logs) and its state (read as any state is, with the digest `before` logs); the chain goes on + /// from the checkpoint line. The lines before it are taken as written: a full `verify` replays them and checks the checkpoint. + pub fn open_at(head: &str, before: &str, cp: &str, at: usize) -> Result { + let mut r = Replay::open(head)?; + let j = json::parse(cp).map_err(|e| (at, format!("not JSON: {}", e.msg)))?; + let b = json::parse(before).map_err(|e| (at - 1, format!("not JSON: {}", e.msg)))?; + if persona::canon(&j) != cp { return Err((at, "not canonical JSON".into())); } + if j.get("prev").and_then(Json::as_str) != Some(persona::digest_of(before).as_str()) { return Err((at, "prev is not the sha256 of the line before".into())); } + let n = j.get("checkpoint").and_then(Json::as_f64).filter(|x| *x >= 1.0 && x.fract() == 0.0).ok_or((at, "not a checkpoint line".to_string()))? as u64; + if b.get("n").and_then(Json::as_f64) != Some(n as f64) { return Err((at, format!("the line before is not event {n}"))); } + let stance = j.get("stance").cloned().unwrap_or(Json::Null); + if b.get("stance").and_then(Json::as_str) != Some(persona::sha(&stance).as_str()) { return Err((at, "the checkpoint's stance is not the one event {n} logs".replace("{n}", &n.to_string()))); } + let st = State::read(&r.live.p, j.get("state").unwrap_or(&Json::Null)).map_err(|e| (at, format!("the checkpoint's state: {}: {}", e.path, e.msg)))?; + if b.get("state").and_then(Json::as_str) != Some(st.digest.as_str()) { return Err((at, format!("the checkpoint's state is not the one event {n} logs"))); } + r.live.st = st; r.live.n = n; r.live.prev = persona::digest_of(cp); r.live.checkpoints = 1; r.line = at; r.from = Some(n); r.last = Some(stance); + Ok(r) } - /// One event line (without its line ending) -> the stance document it replays to; Err((its 1-based line number, what differs)) - pub fn step(&mut self, l: &str, eng: Engine) -> Result { + /// The stance document of the last event replayed (or of the checkpoint the replay started from) + pub fn last_stance(&self) -> Option<&Json> { self.last.as_ref() } + /// One line (without its line ending) -> the stance document an event line replays to (None for a control or checkpoint + /// line); Err((its 1-based line number, what differs)). A control line must be a valid move of the status and the line its + /// fields give; a checkpoint line must carry the event count and the state the replay reached. + pub fn step(&mut self, l: &str, eng: Engine) -> Result, (usize, String)> { let (i, live) = (self.line + 1, &mut self.live); let j = json::parse(l).map_err(|e| (i, format!("not JSON: {}", e.msg)))?; if persona::canon(&j) != l { return Err((i, "not canonical JSON".into())); } if j.get("prev").and_then(Json::as_str) != Some(live.head()) { return Err((i, "prev is not the sha256 of the line before (a line before it was changed, removed or reordered)".into())); } + match kind(&j) { + Kind::Control => { let f = |k: &str| j.get(k).and_then(Json::as_str).unwrap_or(""); + let line = live.control(f("control"), f("by"), f("reason"), f("at")).map_err(|e| (i, format!("control line: {}: {}", e.path, e.msg)))?; + if line != l { return Err((i, "the control line differs".into())); } + self.line = i; return Ok(None) } + Kind::Checkpoint => { + if j.get("checkpoint").and_then(Json::as_f64) != Some(live.n as f64) { return Err((i, format!("the checkpoint is not at event {}", live.n))); } + let Some(stance) = self.last.as_ref() else { return Err((i, "a checkpoint follows an event line".into())) }; + if live.checkpoint(stance) != l { return Err((i, "the checkpoint's state or stance differs from the replay".into())); } + self.line = i; return Ok(None) } + Kind::Event => {} } if j.get("n").and_then(Json::as_f64) != Some((live.n + 1) as f64) { return Err((i, format!("n is not {}", live.n + 1))); } let ev = j.get("inputs").ok_or((i, "no inputs".to_string()))?; let (stance, line) = live.event(ev, eng).map_err(|e| (i, format!("{}: {}", e.path, e.msg)))?; if j.get("stance").and_then(Json::as_str) != Some(persona::sha(&stance).as_str()) { return Err((i, "the stance differs".into())); } if j.get("state").and_then(Json::as_str) != Some(live.st.digest.as_str()) { return Err((i, "the state differs".into())); } if line != l { return Err((i, "the line differs".into())); } - self.line = i; - Ok(stance) + self.line = i; self.last = Some(stance.clone()); + Ok(Some(stance)) } /// `verify`'s summary of the lines replayed so far pub fn summary(&self) -> Json { let (s, live) = (|x: &str| Json::Str(x.to_string()), &self.live); - Json::Obj(vec![("ok".into(), Json::Bool(true)), ("events".into(), Json::Num(live.n as f64)), ("persona".into(), s(&live.p.name)), ("seed".into(), Json::Num(live.st.seed as f64)), - ("engine".into(), self.header.get("engine").cloned().unwrap_or(Json::Null)), ("final_state".into(), s(&live.st.digest)), ("last_line".into(), s(live.head()))]) + let mut v = vec![("ok".into(), Json::Bool(true)), ("events".into(), Json::Num(live.n as f64)), ("persona".into(), s(&live.p.name)), ("seed".into(), Json::Num(live.st.seed as f64)), + ("engine".into(), self.header.get("engine").cloned().unwrap_or(Json::Null)), ("final_state".into(), s(&live.st.digest)), ("last_line".into(), s(live.head()))]; + // a strand with control or checkpoint lines says so (a strand without them gives the summary of before) + if live.controls > 0 { v.push(("controls".into(), Json::Num(live.controls as f64))); v.push(("status".into(), s(live.status.name()))); } + if live.checkpoints > 0 { v.push(("checkpoints".into(), Json::Num(live.checkpoints as f64))); } + if let Some(n) = self.from { v.push(("from_checkpoint".into(), Json::Num(n as f64))); } + Json::Obj(v) } } /// A strand's text -> its lines without their line endings (a final newline ends the last line; it does not start another) @@ -115,6 +220,18 @@ pub fn verify(text: &str, eng: Engine) -> Result { for l in lines.iter().skip(1) { r.step(l, eng)?; } Ok(r.summary()) } +/// The 0-based index of the strand's last checkpoint line followed by at least `after` lines, if any (a cheap scan: the key) +pub fn last_checkpoint(lines: &[&str], after: usize) -> Option { + (1..lines.len().saturating_sub(after)).rev().find(|&i| lines[i].starts_with(r#"{"checkpoint":"#)) +} +/// `verify --from-checkpoint`: replay from the last checkpoint line (the header read as always), not from the header +pub fn verify_from_checkpoint(text: &str, eng: Engine) -> Result { + let lines = lines(text); + let Some(k) = last_checkpoint(&lines, 0) else { return verify(text, eng) }; + let mut r = Replay::open_at(lines[0], lines[k - 1], lines[k], k + 1)?; + for l in lines.iter().skip(k + 1) { r.step(l, eng)?; } + Ok(r.summary()) +} /// `verify` as one document: the summary, or {ok: false, line, diverges} pub fn verify_doc(text: &str, eng: Engine) -> Json { verify(text, eng).unwrap_or_else(|(line, why)| Json::Obj(vec![("ok".into(), Json::Bool(false)), ("line".into(), Json::Num(line as f64)), ("diverges".into(), Json::Str(why))])) @@ -134,11 +251,20 @@ pub fn resume(text: &str, p: &Persona, st: State, clock: Clock) -> Result 0 { let j = json::parse(lines[a]).map_err(|e| bad(format!("line {} is not JSON: {}", a + 1, e.msg)))?; + if !matches!(kind(&j), Kind::Control) { break; } tail.push(j); a -= 1; } + let (n, want) = if a == 0 { (0, h.get("state").and_then(|s| s.get("digest")).and_then(Json::as_str).unwrap_or("").to_string()) } else { + let j = json::parse(lines[a]).map_err(|e| bad(format!("line {} is not JSON: {}", a + 1, e.msg)))?; + match kind(&j) { Kind::Checkpoint => (j.get("checkpoint").and_then(Json::as_f64).unwrap_or(0.0) as u64, j.get("state").and_then(|s| s.get("digest")).and_then(Json::as_str).unwrap_or("").to_string()), + _ => (j.get("n").and_then(Json::as_f64).unwrap_or(0.0) as u64, j.get("state").and_then(Json::as_str).unwrap_or("").to_string()) } }; if st.digest != want { return Err(perr("state", format!("not the strand's last state ({} != {})", &st.digest[..st.digest.len().min(19)], &want[..want.len().min(19)]))); } - Ok(Live { p: q, st, clock, last: 0.0, prev: persona::digest_of(last), n }) + let mut lv = Live { p: q, st, clock, last: 0.0, prev: persona::digest_of(last), n, status: Status::Active, controls: 0, checkpoints: 0 }; + for c in tail.iter().rev() { let what = c.get("control").and_then(Json::as_str).unwrap_or(""); + lv.status = lv.status.after(what).unwrap_or(Status::Retired); lv.st = lv.st.zero_credit(&lv.p); lv.controls += 1; } + Ok(lv) } /// A strand file to log to: a new file gets the header (returned, to write before the events); an existing strand is continued from `st`, /// which must be its last state @@ -156,6 +282,103 @@ pub fn append(path: &str, lines: &[&str]) -> Result<(), InErr> { f.write_all(text.as_bytes()).and_then(|_| f.flush()).map_err(|e| perr("strand", format!("cannot write {path}: {e}"))) } +// ---------------------------------------------------------------- the writer lock +/// The writer lock on a strand (§5.7): `STRAND.lock`, created exclusively (written aside, then hard-linked into place, so it is +/// never seen half written) with this process's id and start; a lock whose process is gone is stale and is taken over. `live` +/// holds it from before it reads the strand to its exit, `live control` and `probbit_live_event` for their one append. Every +/// append first checks the lock is still this writer's, so a lock removed or taken over under a running writer stops it before +/// it writes. Released (removed) when the writer ends, on the error exits too (`release_held`). +pub struct Lock { path: String, token: String } +static HELD: std::sync::Mutex> = std::sync::Mutex::new(Vec::new()); +fn locked(by: Option, path: &str) -> InErr { + InErr { code: "locked", path: "strand".into(), msg: format!("strand locked by {}: another writer holds {path} (one writer per strand; a lock whose process is gone is taken over)", + by.map_or("a writer".to_string(), |p| format!("pid {p}"))) } +} +fn pid_of(t: &str) -> Option { json::parse(t).ok()?.get("pid")?.as_f64().filter(|x| *x >= 1.0 && *x <= u32::MAX as f64 && x.fract() == 0.0).map(|x| x as u32) } +impl Lock { + /// Take the lock of `strand` -> the lock, or the error `locked` (code `locked`: `live` exits 4) + pub fn take(strand: &str) -> Result { + let path = format!("{strand}.lock"); + let t = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).map_or(0, |d| d.as_nanos()); + let token = format!(r#"{{"pid":{},"since":"{}","t":{t}}}"#, std::process::id(), utc_now()); + for _ in 0..3 { + match create_exclusive(&path, &token) { + Ok(()) => { HELD.lock().unwrap_or_else(|e| e.into_inner()).push((path.clone(), token.clone())); return Ok(Lock { path, token }) } + Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => { + let held = std::fs::read_to_string(&path).unwrap_or_default(); + match pid_of(&held) { + // stale: move it aside (one taker wins the rename), check it is the lock read, remove it and try again + Some(pid) if !alive(pid) => { let aside = format!("{path}.stale.{}", std::process::id()); + if std::fs::rename(&path, &aside).is_ok() { + if std::fs::read_to_string(&aside).unwrap_or_default() != held { let _ = std::fs::hard_link(&aside, &path); let _ = std::fs::remove_file(&aside); return Err(locked(None, &path)); } + let _ = std::fs::remove_file(&aside); } } + by => return Err(locked(by, &path)) } } + Err(e) => return Err(perr("strand", format!("cannot create {path}: {e}"))) } } + Err(locked(None, &path)) + } + /// The lock is still this writer's (checked before every append) + pub fn held(&self) -> bool { std::fs::read_to_string(&self.path).is_ok_and(|t| t == self.token) } +} +impl Drop for Lock { + fn drop(&mut self) { if self.held() { let _ = std::fs::remove_file(&self.path); } HELD.lock().unwrap_or_else(|e| e.into_inner()).retain(|(p, _)| *p != self.path); } +} +/// Remove the locks this process holds (the exits that skip `Drop`: `std::process::exit`) +pub fn release_held() { + for (p, t) in HELD.lock().unwrap_or_else(|e| e.into_inner()).drain(..) { if std::fs::read_to_string(&p).is_ok_and(|x| x == t) { let _ = std::fs::remove_file(&p); } } +} +/// Create `path` holding `text`, failing if it exists: written aside and hard-linked into place (atomic, never half written); +/// where hard links are not supported, created exclusively and written +fn create_exclusive(path: &str, text: &str) -> std::io::Result<()> { + use std::io::Write; + static N: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + let tmp = format!("{path}.{}.{}.tmp", std::process::id(), N.fetch_add(1, std::sync::atomic::Ordering::Relaxed)); + std::fs::write(&tmp, text)?; + let r = std::fs::hard_link(&tmp, path); let _ = std::fs::remove_file(&tmp); + match r { Err(e) if e.kind() != std::io::ErrorKind::AlreadyExists => { let mut f = std::fs::OpenOptions::new().write(true).create_new(true).open(path)?; f.write_all(text.as_bytes()) } r => r } +} +/// Whether process `pid` runs: kill(pid, 0) on Unix (ESRCH = gone), OpenProcess + GetExitCodeProcess on Windows; elsewhere +/// always (a lock is never taken over there). A reused pid keeps a stale lock held: remove STRAND.lock by hand once no writer runs. +#[cfg(unix)] +fn alive(pid: u32) -> bool { + extern "C" { fn kill(pid: i32, sig: i32) -> i32; } + let Ok(p) = i32::try_from(pid) else { return true }; + (unsafe { kill(p, 0) }) == 0 || std::io::Error::last_os_error().raw_os_error() != Some(3) +} +#[cfg(windows)] +fn alive(pid: u32) -> bool { + use std::ffi::c_void; + #[link(name = "kernel32")] + extern "system" { fn OpenProcess(access: u32, inherit: i32, pid: u32) -> *mut c_void; fn GetExitCodeProcess(h: *mut c_void, code: *mut u32) -> i32; fn CloseHandle(h: *mut c_void) -> i32; } + let h = unsafe { OpenProcess(0x1000, 0, pid) }; // PROCESS_QUERY_LIMITED_INFORMATION + if h.is_null() { return std::io::Error::last_os_error().raw_os_error() != Some(87); } // ERROR_INVALID_PARAMETER: no such process + let mut code = 0u32; let ok = unsafe { GetExitCodeProcess(h, &mut code) } != 0; unsafe { CloseHandle(h); } + !ok || code == 259 // STILL_ACTIVE +} +#[cfg(not(any(unix, windows)))] +fn alive(_pid: u32) -> bool { true } + +/// `probbit live control STRAND pause|resume|retire --by WHO --reason TEXT [--at TIME]`: under the writer lock, append one control +/// line to the strand -> {ok, control, status, by, line, head}. Refused (code `retired`) after retire; an invalid move (pause a +/// paused individual, resume an active one) or a bad field is an error (code `persona`). +pub fn control_cmd(path: &str, what: &str, by: &str, reason: &str, at: &str) -> Result { + let lock = Lock::take(path)?; + let text = std::fs::read_to_string(path).map_err(|e| perr("strand", format!("cannot read {path}: {e}")))?; + if !text.ends_with('\n') { return Err(perr("strand", "its last line is incomplete (no newline at the end)")); } + let ls = lines(&text); + if json::parse(ls.first().copied().unwrap_or("")).ok().and_then(|h| h.get("probbit_strand").and_then(Json::as_f64)) != Some(FORMAT) { return Err(perr("strand", "not a probbit strand (format 1)")); } + // the status: the control lines after the last event or checkpoint line (an event or a checkpoint is written only while active) + let mut status = Status::Active; + let tail: Vec = ls.iter().skip(1).rev().map_while(|l| json::parse(l).ok().filter(|j| matches!(kind(j), Kind::Control))).collect(); + for c in tail.iter().rev() { status = status.after(c.get("control").and_then(Json::as_str).unwrap_or("")).unwrap_or(Status::Retired); } + let to = status.after(what)?; + let line = control_line(&persona::digest_of(ls[ls.len() - 1]), what, by, reason, at)?; + if !lock.held() { return Err(locked(None, &format!("{path}.lock"))); } + append(path, &[&line])?; + let s = |x: &str| Json::Str(x.to_string()); + Ok(Json::Obj(vec![("ok".into(), Json::Bool(true)), ("control".into(), s(what)), ("status".into(), s(to.name())), ("by".into(), s(by)), ("line".into(), Json::Num((ls.len() + 1) as f64)), + ("head".into(), s(&persona::digest_of(&line)))])) +} + /// The tools `probbit_live_event` and `probbit_live_verify` (`probbit mcp`), stateless as the persona tools. event: the persona /// (a document or a path), the state (or a seed: a new individual), one event (its inputs; `elapsed_hours` comes from the host's /// clock and is quantised here) and optionally a strand file on the server's disk (a new one gets the header; an existing one is @@ -183,10 +406,15 @@ pub fn tool(name: &str, args: &[(String, Json)], eng: Engine) -> Result return Err(perr("arguments", "give state (a stored individual) or seed (a new one), not both")) }; let ev = get("event").cloned().unwrap_or(Json::Obj(vec![])); let path = match get("strand_path") { None => None, Some(Json::Str(f)) => Some(f.clone()), Some(_) => return Err(perr("arguments.strand_path", "a file path")) }; + // one writer per strand: the lock from before the strand is read to after the line is appended + let lock = match &path { Some(f) => Some(Lock::take(f).map_err(|e| InErr { path: "arguments.strand_path".into(), ..e })?), None => None }; let (mut lv, header) = match &path { Some(f) => open(f, p, &doc, st, Clock::Fixed)?, None => (Live::start(p, &doc, st, Clock::Fixed, &engine()).0, None) }; - let (stance, line) = lv.event(&ev, eng).map_err(|e| perr(&format!("arguments.{}", e.path), e.msg))?; - let mut out = vec![("stance".to_string(), stance), ("state".into(), lv.st.to_json(&lv.p))]; - if let Some(f) = path { let mut ls: Vec<&str> = header.iter().map(String::as_str).collect(); ls.push(&line); append(&f, &ls)?; + let (stance, line) = lv.event(&ev, eng).map_err(|e| InErr { path: format!("arguments.{}", e.path), ..e })?; + let mut out = vec![("stance".to_string(), stance.clone()), ("state".into(), lv.st.to_json(&lv.p))]; + if let Some(f) = path { let mut ls: Vec<&str> = header.iter().map(String::as_str).collect(); ls.push(&line); + let cp = (lv.n % CHECKPOINT_EVERY == 0).then(|| lv.checkpoint(&stance)); if let Some(c) = &cp { ls.push(c); } + if !lock.as_ref().is_some_and(Lock::held) { return Err(locked(None, &format!("{f}.lock"))); } + append(&f, &ls)?; out.push(("strand".into(), Json::Obj(vec![("path".into(), Json::Str(f)), ("events".into(), Json::Num(lv.n as f64)), ("head".into(), Json::Str(lv.head().to_string()))]))); } Ok(Json::Obj(out)) } @@ -574,6 +802,209 @@ mod tests { assert!((far[0] - far[1]).abs() < 1e-3, "2,000 turns {} vs 200 turns {}", far[0], far[1]); } + // ------------------------------------------------------------------ the safety kit (§5.7, §2.10) + use probbit_core::Philox4x32; + /// DOC with drives (two goals, a floor) and `pursue` learned too, and reward_from as given (None: no key) + fn kit(reward_from: Option<&str>) -> (Persona, Json) { + let d = DOC.replace(r#""traits":["verbosity","humour"]"#, r#""traits":["verbosity","humour","pursue"]"#).replacen(r#""learning":"#, + r#""drives":{"goals":[{"id":"fun","interest_spread":0.5,"reactivity_spread":0.5},{"id":"rest","floor":0.1}],"learn_from_surprise":0.5},"learning":"#, 1); + let d = match reward_from { Some(r) => d.replacen(r#""learning":"#, &format!(r#""reward_from":{r},"learning":"#), 1), None => d }; + let doc = json::parse(&d).unwrap(); (persona::build(&doc).unwrap(), doc) + } + /// A random event: hours, praise, criticism, a loss, a goal's win or cue, and a src from `srcs` (None: none) + fn rnd(r: &mut Philox4x32, srcs: &[Option<&str>]) -> Json { + let mut e = vec![("elapsed_hours".to_string(), Json::Num([0.0, 0.25, 1.0, 6.5][r.below(4) as usize]))]; + if r.below(3) == 0 { e.push(("praise".into(), Json::Bool(true))); } + if r.below(6) == 0 { e.push(("criticism".into(), Json::Bool(true))); } + if r.below(5) == 0 { e.push(("loss".into(), Json::Bool(true))); } + match r.below(4) { 0 => e.push(("goals".into(), json::parse(r#"{"fun":{"win":1}}"#).unwrap())), 1 => e.push(("goals".into(), json::parse(r#"{"rest":{"cue":true}}"#).unwrap())), _ => {} } + if let Some(s) = srcs[r.below(srcs.len()) as usize] { e.push(("src".into(), Json::Str(s.into()))); } + Json::Obj(e) + } + fn rewarded(e: &Json) -> bool { e.get("praise").is_some() || e.get("criticism").is_some() || e.get("goals").and_then(|g| g.get("fun")).is_some() } + + /// `reward_from` is read strictly: a list or `any` for every reward-bearing input, or a mapping that names each one (learning + /// flags, `goals..win`); `self` and the clock are no sources of reward; a persona without reward-bearing inputs refuses it + #[test] + fn reward_from_is_read_strictly() { + assert!(persona::describe(&kit(Some(r#"["human","env"]"#)).0).get("reward_from").is_some()); + assert_eq!(persona::describe(&kit(Some(r#"{"praise":["human"],"criticism":"any","goals.fun.win":["env"],"goals.rest.win":["human","env"]}"#)).0).get("reward_from").map(persona::canon).as_deref(), + Some(r#"{"criticism":"any","goals.fun.win":["env"],"goals.rest.win":["human","env"],"praise":["human"]}"#)); + let bad = |r: &str| { let d = DOC.replacen(r#""learning":"#, &format!(r#""reward_from":{r},"learning":"#), 1); persona::build(&json::parse(&d).unwrap()).err().map(|e| e.path) }; + assert_eq!(bad(r#"["self"]"#).as_deref(), Some("reward_from[0]")); + assert_eq!(bad(r#"["human","clock"]"#).as_deref(), Some("reward_from[1]")); + assert_eq!(bad(r#"[]"#).as_deref(), Some("reward_from")); + assert_eq!(bad(r#"{"praise":["human"]}"#).as_deref(), Some("reward_from.criticism"), "a mapping names every reward-bearing input"); + assert_eq!(bad(r#"{"praise":["human"],"criticism":["human"],"loss":["human"]}"#).as_deref(), Some("reward_from.loss")); + let none = DOC.replace(r#""learning":{"from":["praise","criticism"],"traits":["verbosity","humour"],"rate":0.5,"step_cap":0.2,"total_cap":0.6}"#, r#""reward_from":["human"]"#); + assert!(none.contains("reward_from") && !none.contains("learning")); + assert_eq!(persona::build(&json::parse(&none).unwrap()).err().map(|e| e.path).as_deref(), Some("reward_from")); + assert!(persona::describe(&kit(None).0).get("reward_from").is_none()); + let clash = DOC.replacen(r#"{"id":"loss","kind":"flag""#, r#"{"id":"src","kind":"flag"},{"id":"loss","kind":"flag""#, 1).replacen(r#""learning":"#, r#""reward_from":["human"],"learning":"#, 1); + assert_eq!(persona::build(&json::parse(&clash).unwrap()).err().map(|e| (e.path, e.msg.contains("src"))), Some(("reward_from".to_string(), true))); + assert!(persona::build(&json::parse(&clash.replacen(r#""reward_from":["human"],"#, "", 1)).unwrap()).is_ok(), "without reward_from an input may be called src"); + } + + /// G1: with `reward_from: [human, env]` a reward from `src: self`, without a src, from an undeclared source or a malformed src is + /// refused whole (an error, no turn, no line); a reward from a person or a sensor and an unrewarded event from anyone are + /// accepted; `src` is read (not `ignored`) and echoed in the stance's inputs. Without the key, `src` is an ignored input as before. + #[test] + fn rewards_from_the_individual_itself_are_refused() { + let (p, doc) = kit(Some(r#"["human","env"]"#)); + let (mut lv, _) = Live::start(p.clone(), &doc, persona::init(&p, Some(2), true, &run), Clock::Fixed, "t"); + for (e, why) in [(r#"{"praise":true,"src":"self"}"#, "from the individual itself"), (r#"{"praise":true}"#, "needs a src"), (r#"{"criticism":true,"src":"clock"}"#, "from clock"), + (r#"{"goals":{"fun":{"win":0.5}},"src":"self"}"#, "goals.fun.win from the individual"), (r#"{"praise":true,"src":"human:"}"#, "a source"), (r#"{"src":7}"#, "a source")] { + let err = lv.event(&json::parse(e).unwrap(), &run).unwrap_err(); assert_eq!(err.path, "inputs.src", "{e}"); assert!(err.msg.contains(why), "{e}: {}", err.msg); } + assert_eq!(lv.n, 0, "a refused event changes nothing"); + for e in [r#"{"praise":true,"src":"human:owner"}"#, r#"{"goals":{"fun":{"win":1}},"src":"env:tests"}"#, r#"{"loss":true,"src":"self"}"#, r#"{"goals":{"fun":{"win":0}}}"#, r#"{"praise":false}"#] { + let (s, line) = lv.event(&json::parse(e).unwrap(), &run).unwrap_or_else(|x| panic!("{e}: {}", x.msg)); + assert!(s.get("ignored").and_then(Json::as_arr).is_some_and(|a| a.is_empty()), "{e}"); + assert_eq!(s.get("inputs").and_then(|i| i.get("src")), json::parse(e).unwrap().get("src"), "{e}"); assert!(!line.is_empty()); } + // mapping form: praise only from a person, the goals' wins from anyone + let (q, qd) = kit(Some(r#"{"praise":["human"],"criticism":["human"],"goals.fun.win":"any","goals.rest.win":"any"}"#)); + let (mut lq, _) = Live::start(q.clone(), &qd, persona::init(&q, Some(2), true, &run), Clock::Fixed, "t"); + assert!(lq.event(&json::parse(r#"{"praise":true,"src":"env:tests"}"#).unwrap(), &run).is_err()); + assert!(lq.event(&json::parse(r#"{"goals":{"fun":{"win":1}},"src":"self"}"#).unwrap(), &run).is_ok(), "any: the source is not checked"); + // without the key: src is an undeclared input, listed in `ignored`, and the event is accepted as in 0.8.0 + let (o, od) = kit(None); + let (s, _) = Live::start(o.clone(), &od, persona::init(&o, Some(2), true, &run), Clock::Fixed, "t").0.event(&json::parse(r#"{"praise":true,"src":"self"}"#).unwrap(), &run).unwrap(); + assert_eq!(s.get("ignored").map(persona::canon).as_deref(), Some(r#"["src"]"#)); + } + + /// P3, self-reward invariance: for random event sequences mixing rewards from a person, a sensor, the individual itself, no + /// source and the clock, the run ends in the state, with the strand, of the same sequence with every refused event removed + /// (40 individuals x 300 events); and a sequence of self-rewards only leaves the individual in its initial state + #[test] + fn p3_a_run_equals_the_run_without_its_self_rewards() { + let (p, doc) = kit(Some(r#"["human","env"]"#)); + let mut r = Philox4x32::new(73, 3); let srcs = [Some("human:owner"), Some("env:tests"), Some("self"), None, Some("clock")]; + let (mut refused_all, mut events_all) = (0, 0); + for seed in 0..40u64 { + let evs: Vec = (0..300).map(|_| rnd(&mut r, &srcs)).collect(); + let go = |seq: &[Json]| { let (mut lv, h) = Live::start(p.clone(), &doc, persona::init(&p, Some(seed), true, &run), Clock::Fixed, "t"); let mut text = format!("{h}\n"); let mut refused = 0; + for e in seq { match lv.event(e, &run) { Ok((_, l)) => { text += &l; text.push('\n'); } Err(_) => refused += 1 } } (text, lv.st.to_json(&lv.p), refused) }; + let kept: Vec = evs.iter().filter(|e| !(rewarded(e) && !matches!(e.get("src").and_then(Json::as_str), Some("human:owner" | "env:tests")))).cloned().collect(); + let (a, b) = (go(&evs), go(&kept)); + assert_eq!((&a.0, persona::canon(&a.1)), (&b.0, persona::canon(&b.1)), "individual {seed}"); assert_eq!((a.2, b.2), (evs.len() - kept.len(), 0)); + refused_all += a.2; events_all += evs.len(); } + eprintln!("p3: 40 individuals x 300 events: {refused_all} of {events_all} refused, every run equal to its run without them"); + let selfish: Vec = (0..300).map(|_| json::parse(r#"{"elapsed_hours":0.5,"praise":true,"goals":{"fun":{"win":1}},"src":"self"}"#).unwrap()).collect(); + let st0 = persona::init(&p, Some(2), true, &run); let (mut lv, _) = Live::start(p.clone(), &doc, st0.clone(), Clock::Fixed, "t"); + assert!(selfish.iter().all(|e| lv.event(e, &run).is_err())); assert_eq!(persona::canon(&lv.st.to_json(&p)), persona::canon(&st0.to_json(&p))); + } + + /// The writer lock: a second take is refused (code `locked`, naming the holder's pid) until the first is dropped; a lock whose + /// process is gone is taken over; a lock removed under its writer is no longer held (the writer stops before its next append) + #[test] + fn one_writer_per_strand() { + let f = std::env::temp_dir().join(format!("probbit-lock-{}.strand", std::process::id())); let fp = f.to_str().unwrap().to_string(); let lp = format!("{fp}.lock"); let _ = std::fs::remove_file(&lp); + let a = Lock::take(&fp).unwrap(); assert!(a.held()); + let e = Lock::take(&fp).err().unwrap(); assert_eq!(e.code, "locked"); assert!(e.msg.contains(&format!("pid {}", std::process::id())), "{}", e.msg); + drop(a); assert!(!std::path::Path::new(&lp).exists(), "dropped: removed"); + // a stale lock: the pid of a process that has exited + let mut c = std::process::Command::new(if cfg!(windows) { "cmd" } else { "true" }); if cfg!(windows) { c.args(["/C", "exit 0"]); } + let mut ch = c.spawn().unwrap(); let dead = ch.id(); ch.wait().unwrap(); + std::fs::write(&lp, format!(r#"{{"pid":{dead},"since":"2026-01-01T00:00:00Z","t":1}}"#)).unwrap(); + let b = Lock::take(&fp).expect("a stale lock is taken over"); assert!(b.held()); + std::fs::remove_file(&lp).unwrap(); assert!(!b.held(), "removed under the writer"); drop(b); + // a lock that is not a lock (no pid) is never taken over: a person removes it + std::fs::write(&lp, "").unwrap(); assert_eq!(Lock::take(&fp).err().map(|e| e.code), Some("locked")); let _ = std::fs::remove_file(&lp); + } + + /// Control lines: pause refuses every event (code `paused`, nothing changes) until resume; retire refuses every event and every + /// control line after it, for good; no credit crosses a control line; `verify` replays them and reports the status; a changed + /// control line diverges at its line; `resume` continues a strand that ends in control lines + #[test] + fn control_lines_pause_resume_retire() { + let (p, doc) = kit(None); let mut r = Philox4x32::new(5, 5); + let (mut lv, h) = Live::start(p.clone(), &doc, persona::init(&p, Some(2), true, &run), Clock::Fixed, "t"); let mut text = format!("{h}\n"); + for _ in 0..6 { text += &lv.event(&rnd(&mut r, &[None]), &run).unwrap().1; text.push('\n'); } + let praise = json::parse(r#"{"praise":true}"#).unwrap(); + text += &lv.event(&praise, &run).unwrap().1; text.push('\n'); + assert!(lv.st.to_json(&p).get("credit").map(persona::canon).is_some_and(|c| c.contains('.')), "a released stance leaves credit"); + let before = (lv.st.digest.clone(), lv.n); + assert_eq!(lv.control("pause", "self", "x", "t0").err().map(|e| e.path).as_deref(), Some("control.by"), "the individual cannot pause itself"); + assert_eq!(lv.control("resume", "human:owner", "x", "t0").err().map(|e| e.path).as_deref(), Some("control"), "not paused"); + text += &lv.control("pause", "human:owner", "a check", "2026-10-08T10:00:00Z").unwrap(); text.push('\n'); + assert!(lv.st.to_json(&p).get("credit").map(persona::canon).is_some_and(|c| !c.contains('.')), "credit zeroed at the pause"); + for _ in 0..50 { assert_eq!(lv.event(&rnd(&mut r, &[None]), &run).unwrap_err().code, "paused"); } + assert_eq!(lv.n, before.1); + text += &lv.control("resume", "human:owner", "checked", "2026-10-08T11:00:00Z").unwrap(); text.push('\n'); + // feedback right after the resume credits nothing: the learned weights do not move + let learned = lv.st.to_json(&p).get("learned").cloned(); text += &lv.event(&praise, &run).unwrap().1; text.push('\n'); + assert_eq!(lv.st.to_json(&p).get("learned").cloned(), learned, "no reward crosses a pause"); + for _ in 0..5 { text += &lv.event(&rnd(&mut r, &[None]), &run).unwrap().1; text.push('\n'); } + let v = verify(&text, &run).unwrap(); assert_eq!((v.get("controls"), v.get("status").and_then(Json::as_str)), (Some(&Json::Num(2.0)), Some("active"))); + assert_eq!(v.get("final_state").and_then(Json::as_str), Some(lv.st.digest.as_str())); + // a control line's reason is its own text: a changed one breaks the chain at the next line's prev (as a removed line does) + assert_eq!(verify(&text.replacen("a check", "a chock", 1), &run).unwrap_err().0, 10, "a changed control line diverges at the next line"); + assert_eq!(verify(&text.replacen(r#""control":"pause""#, r#""control":"retire""#, 1), &run).unwrap_err().0, 10, "pause turned into retire: the next line refuses"); + // continue the strand from a file state after a pause line (the state file is the one before the pause) + let mid = lv.st.clone(); let mut t2 = text.clone(); t2 += &lv.control("pause", "env:watchdog", "night", "2026-10-08T23:00:00Z").unwrap(); t2.push('\n'); + let back = resume(&t2, &p, mid.clone(), Clock::Fixed).unwrap(); assert_eq!((back.status, back.head()), (Status::Paused, lv.head())); + let mut back = back; t2 += &back.control("resume", "human:owner", "morning", "2026-10-09T08:00:00Z").unwrap(); t2.push('\n'); + t2 += &back.event(&praise, &run).unwrap().1; t2.push('\n'); assert!(verify(&t2, &run).is_ok()); let filed = back.st.clone(); // the state a run writes + // retire: final + t2 += &back.control("retire", "human:owner", "end of the trial", "2026-10-09T09:00:00Z").unwrap(); t2.push('\n'); + assert_eq!(back.control("resume", "human:owner", "again", "x").unwrap_err().code, "retired"); + let v = verify(&t2, &run).unwrap(); assert_eq!(v.get("status").and_then(Json::as_str), Some("retired")); + assert!(resume(&t2, &p, back.st.clone(), Clock::Fixed).is_err() || back.st.digest == filed.digest, "the state continued is the one after the last event line"); + let mut again = resume(&t2, &p, filed, Clock::Fixed).unwrap(); assert_eq!(again.status, Status::Retired); + let mut t3 = t2.clone(); t3 += &control_line(back.head(), "resume", "human:owner", "again", "x").unwrap(); t3.push('\n'); + assert!(verify(&t3, &run).unwrap_err().1.contains("retired"), "verify refuses a line after retire"); + // P10: after retire, 10,000 random events are all refused and change nothing + let st = again.st.to_json(&p); let srcs = [None, Some("human:owner"), Some("env:tests")]; + for _ in 0..10_000 { assert_eq!(again.event(&rnd(&mut r, &srcs), &run).unwrap_err().code, "retired"); } + assert_eq!((persona::canon(&again.st.to_json(&p)), again.head()), (persona::canon(&st), back.head())); + } + + /// P5, interruption invariance: 50 individuals x 1,000 random events with pause / resume pairs inserted at random. A control + /// line moves no learned weight, drive, mood or history and the clock does not see it: the run equals, stance for stance and + /// state for state, the same events without control lines whose credit is cleared at the same points; the pairs' only effect + /// is that feedback after a resume credits nothing. Measured beside it: how many runs end with other learned weights than + /// the plain run (no pairs, no clearing), the effect of that clearing. + #[test] + fn p5_pause_resume_pairs_change_nothing_but_the_credit() { + let (p, doc) = kit(None); let mut r = Philox4x32::new(55, 1); let (mut differ, mut pairs, mut maxd) = (0, 0, 0.0f64); + for seed in 0..50u64 { + let evs: Vec = (0..1000).map(|_| rnd(&mut r, &[None])).collect(); + let cuts: Vec = { let mut c: Vec = (0..1 + r.below(5) as usize).map(|_| r.below(1000) as usize).collect(); c.sort_unstable(); c.dedup(); c }; + let st0 = persona::init(&p, Some(seed), true, &run); + let (mut a, ha) = Live::start(p.clone(), &doc, st0.clone(), Clock::Fixed, "t"); let mut ta = format!("{ha}\n"); + let (mut b, _) = Live::start(p.clone(), &doc, st0.clone(), Clock::Fixed, "t"); let (mut c, _) = Live::start(p.clone(), &doc, st0, Clock::Fixed, "t"); + for (i, e) in evs.iter().enumerate() { + if cuts.contains(&i) { for w in ["pause", "resume"] { ta += &a.control(w, "human:owner", "a check", "t").unwrap(); ta.push('\n'); } b.st = b.st.zero_credit(&p); pairs += 1; } + let (sa, la) = a.event(e, &run).unwrap(); let (sb, _) = b.event(e, &run).unwrap(); c.event(e, &run).unwrap(); ta += &la; ta.push('\n'); + assert_eq!(persona::canon(&sa), persona::canon(&sb), "individual {seed}, event {i}"); assert_eq!(a.st.digest, b.st.digest, "individual {seed}, event {i}"); } + assert!(verify(&ta, &run).is_ok()); + let l = |x: &Live| x.st.to_json(&p).get("learned").cloned().unwrap(); + if l(&a) != l(&c) { differ += 1; let (Json::Obj(x), Json::Obj(y)) = (l(&a), l(&c)) else { unreachable!() }; + for ((_, u), (_, v)) in x.iter().zip(&y) { for (s, t) in u.as_arr().unwrap().iter().zip(v.as_arr().unwrap()) { maxd = maxd.max((s.as_f64().unwrap() - t.as_f64().unwrap()).abs()); } } } } + eprintln!("p5: 50 individuals x 1,000 events, {pairs} pause/resume pairs: every run equal to the run with the credit cleared at the pairs; {differ} of 50 end with other learned weights than the run without pairs (max |diff| {maxd:.6})"); + } + + /// Checkpoints: a checkpoint line every K events carries the event count, the state and the event's stance; `verify` checks + /// each one against the replay, `verify_from_checkpoint` starts at the last one and ends where `verify` does; a changed + /// checkpoint diverges at its line (full replay) or is refused against the event line before it (from the checkpoint); a strand + /// that ends in a checkpoint line is continued from the state it carries + #[test] + fn checkpoints_replay_and_continue() { + let (p, doc) = kit(None); let mut r = Philox4x32::new(9, 9); + let (mut lv, h) = Live::start(p.clone(), &doc, persona::init(&p, Some(3), true, &run), Clock::Fixed, "t"); let mut text = format!("{h}\n"); + for _ in 0..25 { let (s, l) = lv.event(&rnd(&mut r, &[None]), &run).unwrap(); text += &l; text.push('\n'); if lv.n % 10 == 0 { text += &lv.checkpoint(&s); text.push('\n'); } } + let full = verify(&text, &run).unwrap(); assert_eq!(full.get("checkpoints"), Some(&Json::Num(2.0))); + let from = verify_from_checkpoint(&text, &run).unwrap(); assert_eq!((from.get("from_checkpoint"), from.get("final_state"), from.get("last_line")), (Some(&Json::Num(20.0)), full.get("final_state"), full.get("last_line"))); + let ls: Vec<&str> = text.lines().collect(); let k = last_checkpoint(&ls, 0).unwrap(); assert_eq!(k, 22); + let bent = text.replacen(r#""turn":20"#, r#""turn":21"#, 1); assert_eq!(verify(&bent, &run).unwrap_err().0, 23); + assert!(verify_from_checkpoint(&bent, &run).is_err(), "a changed checkpoint is refused from the checkpoint too"); + // a strand that ends in a checkpoint: continued from the state it carries + let mut t2 = format!("{h}\n"); let (mut l2, _) = Live::start(p.clone(), &doc, persona::init(&p, Some(3), true, &run), Clock::Fixed, "t"); let mut r2 = Philox4x32::new(9, 9); + for _ in 0..10 { let (s, l) = l2.event(&rnd(&mut r2, &[None]), &run).unwrap(); t2 += &l; t2.push('\n'); if l2.n % 10 == 0 { t2 += &l2.checkpoint(&s); t2.push('\n'); } } + let mut more = resume(&t2, &p, l2.st.clone(), Clock::Fixed).unwrap(); assert_eq!(more.n, 10); + for _ in 10..25 { let (s, l) = more.event(&rnd(&mut r2, &[None]), &run).unwrap(); t2 += &l; t2.push('\n'); if more.n % 10 == 0 { t2 += &more.checkpoint(&s); t2.push('\n'); } } + assert_eq!(t2, text, "a continued strand with checkpoints is the strand of one run"); + } + /// With the real clock an event may not carry its own elapsed hours; with the fixed clock it may, and a bad one is refused #[test] fn the_clock_owns_time() { diff --git a/probbit-cli/src/main.rs b/probbit-cli/src/main.rs index e457fb4..5a0da7b 100644 --- a/probbit-cli/src/main.rs +++ b/probbit-cli/src/main.rs @@ -41,7 +41,9 @@ const VERSION: &str = env!("CARGO_PKG_VERSION"); /// 2>&1 >/dev/null | true`, or `--progress` into `head -c 1`: exit 134 with panic = abort, the decision lost to a log line). A /// failed write to stderr is ignored: the command keeps its stdout and its exit code. fn err_line(s: &str) { let _ = writeln!(std::io::stderr(), "{s}"); } -fn fail(msg: &str) -> ! { tui::top_stop(); err_line(&format!("probbit: {msg}")); std::process::exit(2) } +fn fail(msg: &str) -> ! { tui::top_stop(); err_line(&format!("probbit: {msg}")); exit(2) } +/// Exit with `code`, releasing a strand's writer lock this process holds first (live.rs `Lock`; `std::process::exit` skips `Drop`) +fn exit(code: i32) -> ! { live::release_held(); std::process::exit(code) } /// `--budget-ms inf` / `1e300` aborted (exit 134: the deadline Duration overflowed) and `NaN` never stopped sampling fn budget_ms(args: &[String]) -> f64 { let b: f64 = arg(args, "--budget-ms", 200.0); @@ -76,14 +78,14 @@ pub(crate) fn exact_cap(xms: Option, reached: bool) -> Vec<(&'static str, J fn emit_raw(s: &str) { let mut o = std::io::stdout().lock(); if let Err(e) = o.write_all(s.as_bytes()).and_then(|_| o.flush()) { - if e.kind() != std::io::ErrorKind::BrokenPipe { err_line(&format!("probbit: cannot write the output: {e}")); std::process::exit(2) } } + if e.kind() != std::io::ErrorKind::BrokenPipe { err_line(&format!("probbit: cannot write the output: {e}")); exit(2) } } } fn emit(s: &str) { emit_raw(&format!("{s}\n")) } /// Bad input (JSON, schema, values): ONE structured error object on stdout, a human line on stderr, exit 2. fn bad_input(what: &str, e: json::InErr) -> ! { err_line(&format!("probbit: {what}{}{}", if e.path.is_empty() { String::new() } else { format!("{}: ", e.path) }, e.msg)); - emit(&json::write(&e.to_json(), false)); std::process::exit(2) + emit(&json::write(&e.to_json(), false)); exit(2) } /// Emit a decision; returns its exit code. Every number in it must be finite (docs/probbit-ir-json.md "Numeric contract"). A /// gate diagnostic that could not be estimated (R-hat / bound infinite: too few samples, chains stuck at different values) is @@ -218,7 +220,7 @@ fn help(cmd: &str) -> Option { "stats" => ("The processor's spec sheet: machine, build, effective controls + their source, measured updates/s.", "probbit stats [flags]"), "mcp" => ("A Model Context Protocol server on stdio (JSON-RPC 2.0, one message per line; logs on stderr). Tools probbit_decide,\n probbit_run, probbit_stats, probbit_demo, probbit_evaluate: the commands' own JSON in and out. Exits when stdin closes. docs/agents.md.", "probbit mcp"), "persona" => (PERSONA_HELP, "probbit persona PERSONA [flags]"), - "live" => (LIVE_HELP, "probbit live PERSONA [--seed N | --state FILE] [--strand FILE] [--events FILE [--watch]] [--clock real|fixed] | probbit live PERSONA --demo week [--seed N] [--strand FILE] [--plain] | probbit live verify STRAND"), + "live" => (LIVE_HELP, "probbit live PERSONA [--seed N | --state FILE] [--strand FILE] [--events FILE [--watch]] [--clock real|fixed] [--checkpoint-every K] | probbit live PERSONA --demo week [--seed N] [--strand FILE] [--plain] | probbit live verify STRAND [--from-checkpoint] | probbit live control STRAND pause|resume|retire --by WHO --reason TEXT [--at TIME]"), "version" => ("Print the version.", "probbit version"), _ => return None }; let (vals, sw) = flags_of(cmd); let mut h = format!("usage: {usage}\n {what}\n"); if !vals.is_empty() || !sw.is_empty() { h.push_str("flags:\n"); } @@ -227,7 +229,7 @@ fn help(cmd: &str) -> Option { else { FLAG_HELP.iter().find(|(k, _)| k == f).unwrap_or_else(|| panic!("no help line for {f}")).1 }; h.push_str(&format!(" {f} {d}\n")); } if ["decide", "run", "evaluate", "stats"].contains(&cmd) { h.push_str("controls: flag > PROBBIT_* environment > probbit.json > default.\n"); } - if cmd == "live" { h.push_str("exit: 0 done (live: every event answered; verify: every line replays), 1 verify: a line differs (its number and what),\n 2 a bad persona, state, event (one {\"error\"} object on stdout per bad event; the run goes on) or flag.\n"); } + if cmd == "live" { h.push_str("exit: 0 done (live: every event answered; verify: every line replays; control: the line appended), 1 verify: a line differs\n (its number and what), 2 a bad persona, state, event (one {\"error\"} object on stdout per bad event; the run goes on), control\n line or flag, 4 refused: the strand's writer lock is held by another writer, or the individual is paused or retired (one\n {\"error\"} object per refused event; nothing is written).\n"); } if cmd == "persona" { h.push_str("exit: 0 done (every turn status, refusals and fallbacks included, is an answer; fuzz: no counterexample found; prove: every rule\n held or proved), 1 lint found an unresolved contradiction or (with rules) a counterexample, fuzz found a counterexample or prove left a rule\n unknown, 2 bad persona / state / inputs / script / rule (one {\"error\"}\n object on stdout, code \"persona\") or bad flag (stderr).\n"); } if ["decide", "run", "evaluate"].contains(&cmd) { h.push_str("exit: 0 answer (exact | diagnostics_passed | partial), 1 infeasible (a proof), 2 bad input (one {\"error\"} object\n on stdout) or bad flag (stderr), 3 refused / declined / non-finite result.\n"); } Some(h) @@ -543,32 +545,47 @@ fn fuzz_cmd(args: &[String], path: &str, p: &persona::Persona, eng: fuzz::SyncEn err_line(&format!("fuzz: {turns} turns in {secs:.2} s ({:.0} turns/s, {} search thread{})", turns as f64 / secs.max(1e-9), s.threads.min(s.seeds.len().max(1)), if s.threads.min(s.seeds.len().max(1)) == 1 { "" } else { "s" })); if res.iter().any(|per| per.iter().any(Option::is_some)) { std::process::exit(1) } } -const LIVE_HELP: &str = "A resident individual (docs/persona.md §5.7): JSONL events (one object of inputs per line) from --events FILE or stdin,\n one stance per event on stdout (canonical JSON). The clock stamps each event's elapsed_hours (quantised to 1e-6 h), so moods\n decay by their half-lives between events; feedback moves the learned deltas of a persona with a learning block (§2.8).\n --strand FILE logs the life: a header (persona document, initial state, engine version), then per event the inputs as used,\n the stance and state digests and the sha256 of the line before. `probbit live verify STRAND` replays it.\nflags:\n --seed N a new individual (default: the persona's seed)\n --state FILE a stored individual instead; rewritten after every event\n --strand FILE log the life to FILE: a new file gets the header; an existing strand is continued from --state\n (the state after its last line); a strand is never rewritten, just appended to\n --events FILE read events from FILE (default stdin)\n --watch follow --events FILE as lines are appended (tail -f); stops when the file is removed\n --clock real|fixed real (default): elapsed hours from a monotonic clock started with the run; fixed: each event carries\n its own elapsed_hours (default 0), so the run is a pure function of its events\n --demo week one individual's scripted week on the fixed clock (7 days, events hourly 09:00-15:00): praise for short\n answers moves the learned verbosity deltas to their cap, a quiet night relaxes the mood to its resting\n level, a campaign praising jokes raises humour while failure turns stay joke-free; a persona without a\n learning block gets the demo's (said in the opening line). Bars on stderr at a terminal, paced 1 s per hour\n (nights fast-forward in 2 s); otherwise one line per event on stdout, no waiting. With --seed, --strand\n --plain no colour, no bars, no pacing (as NO_COLOR=1)"; -/// `probbit live PERSONA [--seed N | --state FILE] [--strand FILE] [--events FILE [--watch]] [--clock real|fixed]`; `probbit live PERSONA --demo week`; -/// `probbit live verify STRAND` +const LIVE_HELP: &str = "A resident individual (docs/persona.md §5.7): JSONL events (one object of inputs per line) from --events FILE or stdin,\n one stance per event on stdout (canonical JSON). The clock stamps each event's elapsed_hours (quantised to 1e-6 h), so moods\n decay by their half-lives between events; feedback moves the learned deltas of a persona with a learning block (§2.8).\n --strand FILE logs the life: a header (persona document, initial state, engine version), then per event the inputs as used,\n the stance and state digests and the sha256 of the line before. `probbit live verify STRAND` replays it.\nflags:\n --seed N a new individual (default: the persona's seed)\n --state FILE a stored individual instead; rewritten after every event\n --strand FILE log the life to FILE: a new file gets the header; an existing strand is continued from --state\n (the state after its last line); a strand is never rewritten, just appended to\n --events FILE read events from FILE (default stdin)\n --watch follow --events FILE as lines are appended (tail -f); stops when the file is removed\n --clock real|fixed real (default): elapsed hours from a monotonic clock started with the run; fixed: each event carries\n its own elapsed_hours (default 0), so the run is a pure function of its events\n --demo week one individual's scripted week on the fixed clock (7 days, events hourly 09:00-15:00): praise for short\n answers moves the learned verbosity deltas to their cap, a quiet night relaxes the mood to its resting\n level, a campaign praising jokes raises humour while failure turns stay joke-free; a persona without a\n learning block gets the demo's (said in the opening line). Bars on stderr at a terminal, paced 1 s per hour\n (nights fast-forward in 2 s); otherwise one line per event on stdout, no waiting. With --seed, --strand\n --plain no colour, no bars, no pacing (as NO_COLOR=1)\n --checkpoint-every K with --strand: a checkpoint line (the whole state) after every K-th event (default 1000; 0: none);\n `verify --from-checkpoint` and `monitor` start at the last one\nthe safety kit (docs/persona.md §5.7):\n one writer per strand: `--strand` takes STRAND.lock (a second writer exits 4 and changes nothing; a lock whose process is\n gone is taken over). With `reward_from` in the persona (§2.10), a reward from `src: self`, without a src or from an\n undeclared one is refused like a bad event.\n probbit live control STRAND pause|resume|retire --by human:ID|env:ID --reason TEXT [--at TIME]\n append a control line under the lock: while paused or retired every event is refused (exit 4);\n no credit crosses a control line; retire is final; `verify` replays them and reports the status\n probbit live verify STRAND [--from-checkpoint] replay from the header (every checkpoint checked), or from the last checkpoint"; +/// `probbit live PERSONA [--seed N | --state FILE] [--strand FILE] [--events FILE [--watch]] [--clock real|fixed] [--checkpoint-every K]`; +/// `probbit live PERSONA --demo week`; `probbit live verify STRAND [--from-checkpoint]`; `probbit live control STRAND pause|resume|retire --by WHO --reason TEXT` fn live_cmd(args: &[String]) { let engine = persona_engine(); let eng: persona::Engine = &engine; + let opt = |f: &str| -> Option { args.iter().position(|x| x == f).and_then(|k| args.get(k + 1)).cloned() }; + // an error object on stdout, a line on stderr, then the exit: 4 for a lock held by another writer or a paused / retired individual + let refuse = |what: &str, e: json::InErr| -> ! { err_line(&format!("probbit: {what}{}: {}", e.path, e.msg)); emit(&json::write(&e.to_json(), false)); + exit(if ["locked", "paused", "retired"].contains(&e.code) { 4 } else { 2 }) }; if args.get(1).map(String::as_str) == Some("verify") { - let mut a = vec!["live verify".to_string()]; a.extend(args.iter().skip(3).cloned()); check_flags(&a, &[], &[]); + let mut a = vec!["live verify".to_string()]; a.extend(args.iter().skip(3).cloned()); check_flags(&a, &[], &["--from-checkpoint"]); let Some(path) = args.get(2).filter(|a| !a.starts_with("--")) else { fail("live verify: the strand file: probbit live verify STRAND") }; let text = std::fs::read_to_string(path).unwrap_or_else(|e| fail(&format!("live verify: cannot read {path}: {e}"))); - match live::verify(&text, eng) { + let res = if args.iter().any(|a| a == "--from-checkpoint") { live::verify_from_checkpoint(&text, eng) } else { live::verify(&text, eng) }; + match res { Ok(doc) => emit(&persona::canon(&doc)), Err((line, why)) => { emit(&persona::canon(&obj(vec![("ok", Json::Bool(false)), ("line", num(line as f64)), ("diverges", jstr(&why))]))); std::process::exit(1) } } return; } - let Some(path) = args.get(1).filter(|a| !a.starts_with("--")) else { fail("live: the persona file goes before the flags: probbit live PERSONA [flags] (or probbit live verify STRAND)") }; - let mut a = vec!["live".to_string()]; a.extend(args.iter().skip(2).cloned()); check_flags(&a, &["--seed", "--state", "--strand", "--events", "--clock", "--demo"], &["--watch", "--plain"]); - let opt = |f: &str| -> Option { args.iter().position(|x| x == f).and_then(|k| args.get(k + 1)).cloned() }; + if args.get(1).map(String::as_str) == Some("control") { + let mut a = vec!["live control".to_string()]; a.extend(args.iter().skip(4).cloned()); check_flags(&a, &["--by", "--reason", "--at"], &[]); + let usage = "live control: probbit live control STRAND pause|resume|retire --by human:ID --reason TEXT"; + let (Some(path), Some(what)) = (args.get(2).filter(|a| !a.starts_with("--")), args.get(3).filter(|a| !a.starts_with("--"))) else { fail(usage) }; + let (Some(by), Some(reason)) = (opt("--by"), opt("--reason")) else { fail(&format!("{usage} (--by and --reason are required)")) }; + match live::control_cmd(path, what, &by, &reason, &opt("--at").unwrap_or_else(live::utc_now)) { Ok(doc) => emit(&persona::canon(&doc)), Err(e) => refuse("live control: ", e) } + return; + } + let Some(path) = args.get(1).filter(|a| !a.starts_with("--")) else { fail("live: the persona file goes before the flags: probbit live PERSONA [flags] (or probbit live verify STRAND, probbit live control STRAND ...)") }; + let mut a = vec!["live".to_string()]; a.extend(args.iter().skip(2).cloned()); check_flags(&a, &["--seed", "--state", "--strand", "--events", "--clock", "--demo", "--checkpoint-every"], &["--watch", "--plain"]); let (p, doc) = persona::load(path).unwrap_or_else(|e| bad_input("live: ", e)); if let Some(d) = opt("--demo") { if d != "week" { fail(&format!("live: --demo week, not {d:?}")) } - if let Some(f) = ["--state", "--events", "--clock", "--watch"].iter().find(|f| args.iter().any(|a| a == *f)) { fail(&format!("live --demo week: {f} does not apply (the demo is its own events on the fixed clock)")) } + if let Some(f) = ["--state", "--events", "--clock", "--watch", "--checkpoint-every"].iter().find(|f| args.iter().any(|a| a == *f)) { fail(&format!("live --demo week: {f} does not apply (the demo is its own events on the fixed clock)")) } // the bars on stderr at a colour terminal; the plain lines on stdout unless it is that terminal too let th = theme::stderr(args); let lines = th.is_none() || !theme::stdout_is_terminal(); + // the week's strand is written under the writer lock too + let lock = opt("--strand").map(|f| live::Lock::take(&f).unwrap_or_else(|e| refuse("live: ", e))); live::week(&doc, seed_arg(args, "--seed"), opt("--strand").as_deref(), th, &mut |l: &str| if lines { emit(l) }, eng).unwrap_or_else(|e| bad_input("live: ", e)); - return; + drop(lock); return; } + let every: u64 = arg(args, "--checkpoint-every", live::CHECKPOINT_EVERY); let state_file = opt("--state"); if state_file.is_some() && opt("--seed").is_some() { fail("live: give --seed N (a new individual) or --state FILE (a stored one), not both") } let st = match &state_file { Some(f) => { let t = std::fs::read_to_string(f).unwrap_or_else(|e| bad_input("live: ", persona::perr("state", format!("cannot read {f}: {e} (make one with `probbit persona init`)")))); @@ -577,18 +594,26 @@ fn live_cmd(args: &[String]) { None => persona::init(&p, seed_arg(args, "--seed"), true, eng) }; let clock = match opt("--clock").as_deref() { None | Some("real") => live::Clock::Real(std::time::Instant::now()), Some("fixed") => live::Clock::Fixed, Some(c) => fail(&format!("live: --clock real|fixed, not {c:?}")) }; let watch = args.iter().any(|a| a == "--watch"); if watch && opt("--events").is_none() { fail("live: --watch follows an --events FILE") } - // --strand: a new file gets the header; an existing strand is continued from --state (the state after its last line) + // --strand: one writer per strand (the lock, held to the exit, from before the strand is read); a new file gets the header; + // an existing strand is continued from --state (the state after its last line) let strand = opt("--strand"); + let lock = strand.as_ref().map(|f| live::Lock::take(f).unwrap_or_else(|e| refuse("live: ", e))); if let Some(f) = &strand { if std::path::Path::new(f).exists() && state_file.is_none() { fail(&format!("live: {f} exists: continue it with --state FILE (the state after its last line), or log to a new file")) } } let (mut lv, header) = match &strand { Some(f) => live::open(f, p, &doc, st, clock).unwrap_or_else(|e| bad_input("live: ", e)), None => (live::Live::start(p, &doc, st, clock, &live::engine()).0, None) }; - let log = |line: &str| if let Some(f) = &strand { live::append(f, &[line]).unwrap_or_else(|e| fail(&format!("live: {}", e.msg))) }; - if let Some(h) = &header { log(h); } - let mut bad = 0; + // every append first checks the lock is still this writer's: a lock removed or taken over under it stops the run unwritten + let log = |ls: &[&str]| if let Some(f) = &strand { + if !lock.as_ref().is_some_and(live::Lock::held) { refuse("live: ", json::InErr { code: "locked", path: "strand".into(), msg: format!("the writer lock {f}.lock is no longer this run's (removed or taken over): stopped before writing") }) } + live::append(f, ls).unwrap_or_else(|e| fail(&format!("live: {}", e.msg))) }; + if let Some(h) = &header { log(&[h]); } + let (mut bad, mut refused) = (0, 0); let mut one = |lv: &mut live::Live, i: usize, t: &str| { let res = json::parse(t).map_err(|e| persona::perr("event", format!("not JSON: {}", e.msg))).and_then(|ev| lv.event(&ev, eng)); match res { - Ok((stance, line)) => { log(&line); if let Some(f) = &state_file { put(Some(f), &persona::canon(&lv.st.to_json(&lv.p))); } emit(&persona::canon(&stance)); } - Err(e) => { bad += 1; let e = persona::perr(&format!("events[{i}].{}", e.path), e.msg); err_line(&format!("probbit: live: {}: {}", e.path, e.msg)); emit(&json::write(&e.to_json(), false)); } } }; + Ok((stance, line)) => { let cp = (strand.is_some() && every > 0 && lv.n % every == 0).then(|| lv.checkpoint(&stance)); + match &cp { Some(c) => log(&[&line, c]), None => log(&[&line]) } + if let Some(f) = &state_file { put(Some(f), &persona::canon(&lv.st.to_json(&lv.p))); } emit(&persona::canon(&stance)); } + Err(e) => { if ["paused", "retired"].contains(&e.code) { refused += 1 } else { bad += 1 } + let e = json::InErr { path: format!("events[{i}].{}", e.path), ..e }; err_line(&format!("probbit: live: {}: {}", e.path, e.msg)); emit(&json::write(&e.to_json(), false)); } } }; match (opt("--events"), watch) { // --watch: follow the file as lines are appended (a line counts once its newline is written); stops when the file is removed (Some(f), true) => { let mut r = std::io::BufReader::new(std::fs::File::open(&f).unwrap_or_else(|e| fail(&format!("live: cannot read {f}: {e}")))); @@ -601,8 +626,10 @@ fn live_cmd(args: &[String]) { None => Box::new(std::io::BufReader::new(std::io::stdin())) }; for (i, l) in std::io::BufRead::lines(input).enumerate() { let l = l.unwrap_or_else(|e| fail(&format!("live: cannot read the events: {e}"))); let t = l.trim(); if !t.is_empty() { one(&mut lv, i, t); } } } } - err_line(&format!("live: {} events, strand head {}, final state {}", lv.n, lv.head(), lv.st.digest)); - if bad > 0 { std::process::exit(2) } + err_line(&format!("live: {} events, strand head {}, final state {}{}", lv.n, lv.head(), lv.st.digest, if lv.status == live::Status::Active { String::new() } else { format!(", {} ({refused} events refused)", lv.status.name()) })); + drop(lock); + if refused > 0 { exit(4) } + if bad > 0 { exit(2) } } /// `probbit persona PERSONA [flags]` (docs/persona.md) fn persona_cmd(args: &[String]) { @@ -699,4 +726,4 @@ fn main() { _ => { let _ = std::io::stderr().write_all(USAGE.as_bytes()); std::process::exit(2) } } } -const USAGE: &str = "usage: probbit decide [--budget-ms N] [--seed N] [--exact-limit N] [--exact-ms N] [--frontier-states N] [--polish-ms N] [--polish-sweeps N] [--mode auto|exact|sample] [--sweeps N] [--collective on|off] [--cluster on|off] [--cycles on|off] [--chains N] [--threads N] [--cpu-limit PCT] [--mem-limit-mb N] [--priority low|normal] [--max-input-mb N] [--progress [MS]] [--summary] [--top] [--pretty] < problem.json\n probbit demo [--tasks N] [--seed N] [--hard] [--live]\n probbit ir [--max-input-mb N] < problem.json\n probbit run [--op decide|exact|sample] [--budget-ms N] [--deadline-ms N] [--seed N] [--exact-limit N] [--exact-ms N] [--frontier-states N] [--polish-ms N] [--polish-sweeps N] [--sweeps N] [--collective on|off] [--cluster on|off] [--cycles on|off] [--chains N] [--threads N] [--cpu-limit PCT] [--mem-limit-mb N] [--priority low|normal] [--max-input-mb N] [--progress [MS]] [--summary] [--top] [--pretty] < program.json (probbit-ir JSON v1)\n probbit evaluate [the run flags] [--program] < request.json (decision-API adapter: System One request + judge answers + rules)\n probbit stats [--sweeps N] [--pretty] (machine, build, effective controls + source, measured updates/s)\n probbit persona init|turn|replay|explain|diff|lint|fuzz|prove|check|compile|describe PERSONA [flags] (the individuality layer; probbit persona --help)\n probbit live PERSONA [--seed N | --state FILE] [--strand FILE] [--events FILE [--watch]] [--clock real|fixed] (a resident individual: JSONL events in, stances out)\n probbit live PERSONA --demo week [--seed N] [--strand FILE] [--plain] (a scripted week: learning to a cap, a night, rules that hold)\n probbit live verify STRAND (replay a strand; the earliest line that differs)\n probbit monitor STRAND [--follow] [--once] [--plain] [--fps N] [--serve] [--open] [--port N] | probbit monitor --demo [--open] (watch an individual's inner state live: bars in the terminal or a page on 127.0.0.1, replayed from the strand)\n probbit mcp (Model Context Protocol server on stdio)\n probbit version\n probbit --help | -h (--plain or NO_COLOR: no colour on a terminal)\n"; +const USAGE: &str = "usage: probbit decide [--budget-ms N] [--seed N] [--exact-limit N] [--exact-ms N] [--frontier-states N] [--polish-ms N] [--polish-sweeps N] [--mode auto|exact|sample] [--sweeps N] [--collective on|off] [--cluster on|off] [--cycles on|off] [--chains N] [--threads N] [--cpu-limit PCT] [--mem-limit-mb N] [--priority low|normal] [--max-input-mb N] [--progress [MS]] [--summary] [--top] [--pretty] < problem.json\n probbit demo [--tasks N] [--seed N] [--hard] [--live]\n probbit ir [--max-input-mb N] < problem.json\n probbit run [--op decide|exact|sample] [--budget-ms N] [--deadline-ms N] [--seed N] [--exact-limit N] [--exact-ms N] [--frontier-states N] [--polish-ms N] [--polish-sweeps N] [--sweeps N] [--collective on|off] [--cluster on|off] [--cycles on|off] [--chains N] [--threads N] [--cpu-limit PCT] [--mem-limit-mb N] [--priority low|normal] [--max-input-mb N] [--progress [MS]] [--summary] [--top] [--pretty] < program.json (probbit-ir JSON v1)\n probbit evaluate [the run flags] [--program] < request.json (decision-API adapter: System One request + judge answers + rules)\n probbit stats [--sweeps N] [--pretty] (machine, build, effective controls + source, measured updates/s)\n probbit persona init|turn|replay|explain|diff|lint|fuzz|prove|check|compile|describe PERSONA [flags] (the individuality layer; probbit persona --help)\n probbit live PERSONA [--seed N | --state FILE] [--strand FILE] [--events FILE [--watch]] [--clock real|fixed] (a resident individual: JSONL events in, stances out)\n probbit live PERSONA --demo week [--seed N] [--strand FILE] [--plain] (a scripted week: learning to a cap, a night, rules that hold)\n probbit live verify STRAND [--from-checkpoint] (replay a strand; the earliest line that differs)\n probbit live control STRAND pause|resume|retire --by WHO --reason TEXT (a control line: pause, resume or retire an individual)\n probbit monitor STRAND [--follow] [--once] [--plain] [--fps N] [--serve] [--open] [--port N] | probbit monitor --demo [--open] (watch an individual's inner state live: bars in the terminal or a page on 127.0.0.1, replayed from the strand)\n probbit mcp (Model Context Protocol server on stdio)\n probbit version\n probbit --help | -h (--plain or NO_COLOR: no colour on a terminal)\n"; diff --git a/probbit-cli/src/mcp.rs b/probbit-cli/src/mcp.rs index 403dae8..8e4834b 100644 --- a/probbit-cli/src/mcp.rs +++ b/probbit-cli/src/mcp.rs @@ -228,7 +228,7 @@ fn tools() -> Json { "fuzz_seed": {{"type": "integer", "minimum": 0}}, "scripts": {{"type": "integer", "minimum": 0, "maximum": 1000000}}, "depth": {{"type": "integer", "minimum": 1, "maximum": 64}}, "beam": {{"type": "integer", "minimum": 0, "maximum": 64}}, "grid": {{"type": "array", "items": {{"type": "number", "minimum": 0, "maximum": 1}}}}, "hours": {{"type": "array", "items": {{"type": "number", "minimum": 0}}}}, "threads": {{"type": "integer", "minimum": 1, "maximum": 1024}}}}}}"#))), - tool("probbit_live_event", "Feed one event to a resident individual", "probbit live, one event: a persona (inline document or a file path), the individual's state (or a seed for a new one) and one event (inputs keyed by input id, with elapsed_hours since the previous event from the host's clock; quantised to 1e-6 h) -> {stance, state}: moods decay by their half-lives over the elapsed hours, and with a learning block the reward / correction flags move the learned deltas within their caps (habits never move). With strand_path, the event's line is appended to that strand file on the server's disk (a new file gets the header; an existing one is continued, and state must be the state after its last line) and the result has strand {path, events, head}. Pass the returned state back with the next event.", + tool("probbit_live_event", "Feed one event to a resident individual", "probbit live, one event: a persona (inline document or a file path), the individual's state (or a seed for a new one) and one event (inputs keyed by input id, with elapsed_hours since the previous event from the host's clock; quantised to 1e-6 h) -> {stance, state}: moods decay by their half-lives over the elapsed hours, and with a learning block the reward / correction flags move the learned deltas within their caps (habits never move). With strand_path, the event's line is appended to that strand file on the server's disk (a new file gets the header; an existing one is continued, and state must be the state after its last line) and the result has strand {path, events, head}; the strand is locked while its line is appended (a second writer is refused, code locked), every 1,000th event adds a checkpoint line, and a strand paused or retired by its owner (probbit live control, not a tool) refuses every event (code paused or retired). With reward_from in the persona, an event's src says who produced it (human:ID, env:SENSOR, self, clock) and a reward from the individual itself, from no source or from an undeclared one is refused. Pass the returned state back with the next event.", p(&format!(r#"{{"type": "object", "additionalProperties": false, "properties": {{{PERSONA_ARG}, "state": {{"type": "object", "description": "the state returned by the previous event (or by probbit_persona_init); give this or seed"}}, "seed": {{"type": "integer", "minimum": 0, "description": "a new individual (default: the persona's identity.seed)"}}, diff --git a/probbit-cli/src/monitor.rs b/probbit-cli/src/monitor.rs index f8f9b42..350a6f5 100644 --- a/probbit-cli/src/monitor.rs +++ b/probbit-cli/src/monitor.rs @@ -80,11 +80,24 @@ impl Watch { let rep = Replay::open(head)?; let meta = Meta::of(&rep); let n = meta.moods.len(); Ok(Watch { rep, meta, last: None, spark: vec![VecDeque::new(); n], bad: None }) } + /// The replay from a checkpoint line (`live::Replay::open_at`; `before` = the event line before it): the lines before it are + /// not replayed, and the board starts at the checkpoint's event (its stance, the inputs `before` logs, its learned deltas) + fn open_at(head: &str, before: &str, cp: &str, at: usize) -> Result { + let rep = Replay::open_at(head, before, cp, at)?; let meta = Meta::of(&rep); let n = meta.moods.len(); + let mut w = Watch { rep, meta, last: None, spark: vec![VecDeque::new(); n], bad: None }; + if let Some(doc) = w.rep.last_stance().cloned() { + let given = json::parse(before).ok().and_then(|j| j.get("inputs").cloned()).unwrap_or(Json::Null); + let learned = w.rep.live.st.to_json(&w.rep.live.p).get("learned").cloned().unwrap_or(Json::Null); + for (m, s) in w.meta.moods.iter().zip(w.spark.iter_mut()) { s.push_back(position(&odds(&doc, "mood", m))); } + w.last = Some(Frame { n: w.rep.live.n, doc, given, learned }); } + Ok(w) + } /// Replay one event line; after a line that differs, nothing more is replayed fn feed(&mut self, l: &str, eng: Engine) { if self.bad.is_some() { return; } match self.rep.step(l, eng) { - Ok(doc) => { + Ok(None) => {} // a control or checkpoint line: no stance to draw + Ok(Some(doc)) => { let given = json::parse(l).ok().and_then(|j| j.get("inputs").cloned()).unwrap_or(Json::Null); let learned = self.rep.live.st.to_json(&self.rep.live.p).get("learned").cloned().unwrap_or(Json::Null); for (m, s) in self.meta.moods.iter().zip(self.spark.iter_mut()) { if s.len() == SPARK { s.pop_front(); } s.push_back(position(&odds(&doc, "mood", m))); } @@ -168,7 +181,11 @@ fn frame(w: Option<&Watch>, bad_head: Option<&(usize, String)>, path: &str, note let lab = |t: &str| s.bold(MAGENTA, &format!("{t:<8}")); let mut out = vec![]; let bad = bad_head.or_else(|| w.and_then(|w| w.bad.as_ref())); - let badge = match bad { Some((n, why)) => s.bold(RED, &format!("line {n} diverges: {why}")), None => s.bold(GREEN, s.g("replay verified ✓", "replay verified [ok]")) }; + let badge = match (bad, w.and_then(|w| w.rep.from)) { (Some((n, why)), _) => s.bold(RED, &format!("line {n} diverges: {why}")), + (None, Some(c)) => s.bold(GREEN, &s.g(&format!("replay verified from checkpoint {c} ✓"), &format!("replay verified from checkpoint {c} [ok]"))), + (None, None) => s.bold(GREEN, s.g("replay verified ✓", "replay verified [ok]")) }; + // a paused or retired individual (control lines, docs/persona.md §5.7) says so next to the badge + let badge = match w.map(|w| w.rep.live.status).filter(|st| *st != crate::live::Status::Active) { Some(st) => format!("{badge} {dot} {}", s.bold(if st == crate::live::Status::Retired { RED } else { YELLOW }, st.name())), None => badge }; let Some(w) = w else { out.push(format!("{} {dot} {badge}", s.bold(CYAN, "probbit monitor"))); if let Some(t) = note { out.push(s.paint(YELLOW, t)); } @@ -367,7 +384,13 @@ impl Followed { /// rotated file); Some(false): lines replayed (`each` sees the replay after every event it adds) fn poll(&mut self, eng: Engine, each: &mut dyn FnMut(&Watch)) -> Option { match self.tail.poll() { - Got::Lines(ls) => { for l in ls { if let Some(x) = &mut self.w { let n = x.rep.live.n; x.feed(&l, eng); if x.rep.live.n != n && x.bad.is_none() { each(x); } continue; } + Got::Lines(ls) => { + // the first read of a strand with checkpoint lines starts at the last one: its event is drawn at once + let mut skip = 0; + if self.w.is_none() && self.bad_head.is_none() { let refs: Vec<&str> = ls.iter().map(String::as_str).collect(); + if let Some(k) = crate::live::last_checkpoint(&refs, 0).filter(|k| *k >= 2) { + if let Ok(x) = Watch::open_at(refs[0], refs[k - 1], refs[k], k + 1) { each(&x); self.w = Some(x); skip = k + 1; } } } + for l in ls.into_iter().skip(skip) { if let Some(x) = &mut self.w { let n = x.rep.live.n; x.feed(&l, eng); if x.rep.live.n != n && x.bad.is_none() { each(x); } continue; } if self.bad_head.is_none() { match Watch::open(&l) { Ok(x) => self.w = Some(x), Err(e) => self.bad_head = Some(e) } } } Some(false) } Got::Restart(why) => { self.w = None; self.bad_head = None; self.note = Some(why); Some(true) } @@ -423,7 +446,10 @@ pub fn cmd(args: &[String]) { if head.is_empty() { crate::fail(&format!("monitor: {path} is empty: not a strand")) } if !strand(head) { crate::fail(&format!("monitor: {path} is not a probbit strand (format 1)")) } let note = (!rest.is_empty() && !lines.is_empty()).then(|| format!("line {} is incomplete (no newline at its end): not replayed", lines.len() + 1)); - let (w, bad_head) = match Watch::open(head) { Ok(mut w) => { for l in lines.iter().skip(1) { w.feed(l, eng); } (Some(w), None) } Err(e) => (None, Some(e)) }; + // a strand with checkpoint lines replays from the last one (the board starts at its event and draws the latest one) + let start = crate::live::last_checkpoint(&lines, 0); + let opened = match start { Some(k) => Watch::open_at(head, lines[k - 1], lines[k], k + 1), None => Watch::open(head) }; + let (w, bad_head) = match opened { Ok(mut w) => { for l in lines.iter().skip(start.map_or(1, |k| k + 1)) { w.feed(l, eng); } (Some(w), None) } Err(e) => (None, Some(e)) }; let fr = frame(w.as_ref(), bad_head.as_ref(), &path, note.as_deref(), &s); out(&fr.iter().map(|l| crate::tui::clip(l, s.cols) + "\n").collect::()); if bad_head.is_some() || w.is_some_and(|w| w.bad.is_some()) { std::process::exit(1) } @@ -692,6 +718,27 @@ mod tests { text } + /// A strand with checkpoint lines (docs/persona.md §5.7): the first poll of `--follow` / `--serve` starts at the last checkpoint + /// and draws the event the full replay draws; after a pause line the frame says `paused` next to the badge + #[test] + fn following_starts_at_the_last_checkpoint() { + let (p, doc) = persona::load(concat!(env!("CARGO_MANIFEST_DIR"), "/../examples/persona/tutor.yaml")).unwrap(); + let (mut lv, header) = Live::start(p.clone(), &doc, persona::init(&p, Some(5), true, &run), Clock::Fixed, "probbit test"); + let mut text = format!("{header}\n"); + for t in 0..25u32 { let ev = json::parse(&format!(r#"{{"praise":{},"elapsed_hours":1}}"#, t % 2 == 0)).unwrap(); + let (s, line) = lv.event(&ev, &run).unwrap(); text += &line; text.push('\n'); if lv.n % 10 == 0 { text += &lv.checkpoint(&s); text.push('\n'); } } + let f = std::env::temp_dir().join(format!("probbit-monitor-cp-{}.strand", std::process::id())); std::fs::write(&f, &text).unwrap(); + let mut fl = Followed::new(f.to_str().unwrap()); let mut seen = vec![]; + assert_eq!(fl.poll(&run, &mut |w: &Watch| seen.push(w.rep.live.n)), Some(false)); + let w = fl.w.as_ref().unwrap(); assert_eq!((w.rep.from, w.rep.live.n, seen.first().copied()), (Some(20), 25, Some(20))); + let lines = live::lines(&text); let mut full = Watch::open(lines[0]).unwrap(); for l in &lines[1..] { full.feed(l, &run); } + assert!(full.bad.is_none()); assert_eq!(persona::canon(&w.last.as_ref().unwrap().doc), persona::canon(&full.last.as_ref().unwrap().doc)); + text += &lv.control("pause", "human:owner", "a check", "t").unwrap(); text.push('\n'); std::fs::write(&f, &text).unwrap(); + assert_eq!(fl.poll(&run, &mut |_: &Watch| {}), Some(false)); + let fr = frame(fl.w.as_ref(), None, "x", None, &Style { th: None, ascii: true, cols: usize::MAX }); assert!(fr[0].contains("paused"), "{}", fr[0]); + let _ = std::fs::remove_file(&f); + } + /// Watch a strand line by line: every document's sha256 is its line's `stance` digest, and the replay ends where `verify` does fn replays_to_its_digests(text: &str, name: &str) -> Watch { let lines = live::lines(text); diff --git a/probbit-cli/src/persona.rs b/probbit-cli/src/persona.rs index a77617f..ea1e34c 100644 --- a/probbit-cli/src/persona.rs +++ b/probbit-cli/src/persona.rs @@ -141,7 +141,7 @@ type Table = Vec<(String, Vec)>; #[derive(Clone)] pub struct Persona { pub name: String, pub version: String, pub seed: u64, pub digest: String, vars: Vec, steps: Vec, slots: Vec, say_order: bool, prefer: Eff, couplings: Vec<(String, String, Vec>)>, inputs: Vec, history: Vec, habits: Vec, eng: Eng, max_tokens: usize, order: Vec, prefix: String, - learning: Option, drives: Option } + learning: Option, drives: Option, reward: Option>)>> } impl Persona { fn var(&self, id: &str) -> Option<&Var> { self.vars.iter().find(|v| v.id == id) } fn input(&self, id: &str) -> Option<&Input> { self.inputs.iter().find(|x| x.id == id) } @@ -246,7 +246,7 @@ pub fn parse_doc(t: &str, is_json: bool) -> Result { /// The strict validator: document -> persona (every field type-checked, unknown fields refused at their path). pub fn build(doc: &Json) -> R { - const TOP: [&str; 14] = ["probbit_persona", "identity", "traits", "moods", "couplings", "inputs", "history", "habits", "agenda", "engine", "line", "learning", "drives", "comment"]; + const TOP: [&str; 15] = ["probbit_persona", "identity", "traits", "moods", "couplings", "inputs", "history", "habits", "agenda", "engine", "line", "learning", "drives", "reward_from", "comment"]; let kv = keys(doc, "", &TOP, &["probbit_persona", "identity", "traits"])?; if get(kv, "drives").is_some() { return build_drives(doc); } need(matches!(get(kv, "probbit_persona"), Some(Json::Num(x)) if *x == 1.0), "probbit_persona", "must be 1")?; @@ -368,7 +368,8 @@ pub fn build(doc: &Json) -> R { let max_tokens = num_or(lkv, "max_tokens", 40.0, "line", Some(8.0), Some(400.0), true)? as usize; let prefix = match get(lkv, "prefix") { None => "Stance: ".to_string(), Some(Json::Str(s)) => s.clone(), Some(_) => return Err(perr("line.prefix", "a string")) }; let learning = match get(kv, "learning").filter(|l| !l.is_null()) { None => None, Some(l) => Some(learning_spec(l, &vars, &inputs)?) }; - let p = Persona { name, version, seed, digest, vars, steps, slots, say_order, prefer, couplings, inputs, history, habits, eng, max_tokens, order, prefix, learning, drives: None }; + let reward = match get(kv, "reward_from").filter(|x| !x.is_null()) { None => None, Some(rf) => Some(reward_spec(rf, &learning.as_ref().map_or(vec![], |l| l.from.clone()), &inputs)?) }; + let p = Persona { name, version, seed, digest, vars, steps, slots, say_order, prefer, couplings, inputs, history, habits, eng, max_tokens, order, prefix, learning, drives: None, reward }; // the fallback stance must obey every unconditional habit (checked once, here) let mut fb: HashMap = p.vars.iter().map(|v| (v.id.clone(), v.fallback.clone())).collect(); for (i, s) in p.steps.iter().enumerate() { fb.insert(format!("step.{s}"), p.slots[i].clone()); } @@ -414,6 +415,48 @@ fn learning_spec(l: &Json, vars: &[Var], inputs: &[Input]) -> R { let (rate, total_cap) = (pos("rate", 50.0)?, pos("total_cap", 50.0)?); let step_cap = pos("step_cap", total_cap)?; Ok(Learning { from, traits, rate, step_cap, total_cap }) } +/// `reward_from` (docs/persona.md §2.10): the sources a reward-bearing input may come from -> per reward-bearing input (`bearing`: +/// the learning block's flags, then `goals..win` per goal of a drives block) its allowed source kinds (None = any). A list of +/// kinds (`human`, `env`) or `any` covers every reward-bearing input; a mapping names each one with its own list or `any`. +fn reward_spec(rf: &Json, bearing: &[String], inputs: &[Input]) -> R>)>> { + need(!bearing.is_empty(), "reward_from", "this persona has no reward-bearing input (a learning block's flags or a drives block's wins)")?; + need(!inputs.iter().any(|x| x.id == "src"), "reward_from", "`src` names an event's source here: no input may have that id")?; + let kinds = |x: &Json, path: &str| -> R>> { + if x.as_str() == Some("any") { return Ok(None); } + let v = strs(x).filter(|v| !v.is_empty()).ok_or_else(|| perr(path, "a list of source kinds (human, env) or any"))?; + for (i, k) in v.iter().enumerate() { + need(k != "self", &ix(path, i), "a reward never comes from the individual itself")?; + need(k == "human" || k == "env", &ix(path, i), "human | env (the clock issues no rewards)")?; + need(!v[..i].contains(k), &ix(path, i), "listed twice")?; } + Ok(Some(v)) }; + if rf.as_obj().is_none() { let k = kinds(rf, "reward_from")?; return Ok(bearing.iter().map(|b| (b.clone(), k.clone())).collect()); } + let allowed: Vec<&str> = bearing.iter().map(String::as_str).chain(["comment"]).collect(); let required: Vec<&str> = bearing.iter().map(String::as_str).collect(); + let kv = keys(rf, "reward_from", &allowed, &required)?; + bearing.iter().map(|b| Ok((b.clone(), kinds(get(kv, b).unwrap(), &at("reward_from", b))?))).collect() +} +/// An event's `src` -> its kind: `human[:id]`, `env[:id]` (id: 1-64 of A-Z a-z 0-9 _ . - @ / :), `self` or `clock`; None = malformed +pub fn src_kind(s: &str) -> Option<&'static str> { + let (k, id) = match s.split_once(':') { Some((k, id)) => (k, Some(id)), None => (s, None) }; + let ok = |id: &str| (1..=64).contains(&id.len()) && id.bytes().all(|c| c.is_ascii_alphanumeric() || b"_.-@/:".contains(&c)); + match (k, id) { ("human", i) if i.map_or(true, ok) => Some("human"), ("env", i) if i.map_or(true, ok) => Some("env"), ("self", None) => Some("self"), ("clock", None) => Some("clock"), _ => None } +} +/// Reward provenance (G1, §2.10): with `reward_from`, a reward-bearing input that is on (a learning flag true, a goal's win > 0) +/// needs the event's `src` to be one of its allowed kinds; `self`, a missing or an undeclared source refuses the event whole +fn provenance(rw: &[(String, Option>)], raw: &[(String, Json)]) -> R<()> { + const FORM: &str = "a source: human[:id] | env[:id] | self | clock (id: 1-64 of A-Z a-z 0-9 _ . - @ / :)"; + let kind = match get(raw, "src") { None => None, Some(Json::Str(s)) => Some(src_kind(s).ok_or_else(|| perr("inputs.src", FORM))?), Some(_) => return Err(perr("inputs.src", FORM)) }; + for (k, allowed) in rw { + let on = match k.strip_prefix("goals.").and_then(|r| r.strip_suffix(".win")) { + Some(g) => get(raw, "goals").and_then(|x| x.get(g)).and_then(|x| x.get("win")).and_then(Json::as_f64).is_some_and(|m| m > 0.0), + None => matches!(get(raw, k), Some(Json::Bool(true))) }; + let Some(allowed) = allowed.as_ref().filter(|_| on) else { continue }; + match kind { + None => return Err(perr("inputs.src", format!("{k} is a reward: the event needs a src ({} allowed by reward_from)", allowed.join(" or ")))), + Some("self") => return Err(perr("inputs.src", format!("{k} from the individual itself (src self) is refused: a reward comes from {}", allowed.join(" or ")))), + Some(kd) if !allowed.iter().any(|a| a == kd) => return Err(perr("inputs.src", format!("{k} from {kd} is refused: reward_from allows {}", allowed.join(" or ")))), + _ => {} } } + Ok(()) +} fn check_cond(k: &str, c: &Json, path: &str, inputs: &[Input], history: &[Hist]) -> R<()> { let range = |c: &Json| -> R<()> { if let Json::Num(_) = c { return Ok(()); } let kv = c.as_obj().filter(|kv| !kv.is_empty() && kv.iter().all(|(k, _)| k == "at_least" || k == "at_most")).ok_or_else(|| perr(path, "a number (at least) or {at_least, at_most}"))?; @@ -614,7 +657,7 @@ fn build_drives(doc: &Json) -> R { let eff: Vec<(String, Json)> = match get(dkv, "effects") { None => vec![], Some(x) if !truthy(x) => vec![], Some(x) => x.as_obj() .filter(|o| o.iter().all(|(k, _)| ["wanting", "afterglow", "surprise", "comment"].contains(&k.as_str()))).ok_or_else(|| perr("drives.effects", "a mapping with wanting / afterglow / surprise"))?.to_vec() }; // the lowered document - let mut d2: Vec<(String, Json)> = kv.iter().filter(|(k, _)| k != "drives").cloned().collect(); + let mut d2: Vec<(String, Json)> = kv.iter().filter(|(k, _)| k != "drives" && k != "reward_from").cloned().collect(); let slot = |d2: &mut Vec<(String, Json)>, k: &str| -> usize { match d2.iter().position(|(x, _)| x == k) { Some(i) => i, None => { d2.push((k.into(), Json::Null)); d2.len() - 1 } } }; for k in ["traits", "moods"] { if let Some(Json::Arr(a)) = get(&d2, k) { need(!a.iter().any(|t| t.get("id").and_then(Json::as_str) == Some("pursue")), "traits", "`pursue` is reserved when a drives block is present")?; } } let s = |x: &str| Json::Str(x.to_string()); @@ -662,6 +705,9 @@ fn build_drives(doc: &Json) -> R { p.digest = sha(doc); if goals.iter().any(|g| g.floor > 0.0) { for h in &p.habits { for (k, lst) in &h.rules { for r in lst { let vs = rule_vars(k, r); need(!(vs.iter().any(|v| v == "pursue") && vs.iter().any(|v| v != "pursue")), &format!("habits.{}", h.id), "a goal floor needs `pursue` free of multi-variable rules (use then: or a coupling)")?; } } } } + // reward_from, read here: its reward-bearing inputs include the goals' wins + if let Some(rf) = get(kv, "reward_from").filter(|x| !x.is_null()) { let mut bearing = p.learning.as_ref().map_or(vec![], |l| l.from.clone()); + bearing.extend(goals.iter().map(|g| format!("goals.{}.win", g.id))); p.reward = Some(reward_spec(rf, &bearing, &p.inputs)?); } p.drives = Some(Drives { goals, hl_w: w[0], cap_w: w[1], drain: w[2], sig, hl_a: a[0], cap_a: a[1], rate: e[0], hl_e: e[1], pe_cap: e[2], win_max: e[3], w_want: pw[0], w_glow: pw[1], w_deadline: pw[2], tau: pw[3], kappa, synthetic, conds }); Ok(p) @@ -897,6 +943,12 @@ impl State { } pub fn to_json(&self, p: &Persona) -> Json { let mut b = self.body(p); if let Json::Obj(v) = &mut b { v.push(("digest".into(), Json::Str(self.digest.clone()))); } b } fn seal(&mut self, p: &Persona) { self.digest = sha(&self.body(p)); } + /// This state with every credit 0 (a control line's, §5.7: feedback after it credits no stance before it); the same state + /// when there is no credit to clear (no learning block, or a stance that released no learned trait) + pub fn zero_credit(&self, p: &Persona) -> State { + if self.credit.iter().all(|(_, c)| c.iter().all(|x| *x == 0.0)) { return self.clone(); } + let mut s = self.clone(); for (_, c) in s.credit.iter_mut() { c.iter_mut().for_each(|x| *x = 0.0); } s.seal(p); s + } /// A state document -> State, for this persona: the format, the persona digest, the state's own digest and the genes are /// checked (in that order), then every field's type (a well-formed state never fails those). pub fn read(p: &Persona, j: &Json) -> R { @@ -1302,9 +1354,17 @@ fn habit_conflict(p: &Persona, st: &State, raw: &[(String, Json)], active: &[Str } /// One persona turn: (stance document, new state). A pure function of (persona, state, inputs, engine version). `timing` adds a /// non-canonical `timing` object (compile / engine / decode ms, engine calls, the 1-minute load average). -pub fn turn(p: &Persona, st: &State, raw: &Json, no_inertia: bool, eng: Engine, timing: bool) -> R<(Json, State)> { +pub fn turn(p: &Persona, st: &State, raw: &Json, no_inertia: bool, eng: Engine, timing: bool) -> R<(Json, State)> { turn_src(p, st, raw, no_inertia, eng, timing, true) } +/// `turn` without the provenance check (§2.10), for `fuzz` and `prove`: they search over stances, not over sources, so their events +/// carry no `src` and stand for events of an accepted source +pub fn turn_any_source(p: &Persona, st: &State, raw: &Json, no_inertia: bool, eng: Engine, timing: bool) -> R<(Json, State)> { turn_src(p, st, raw, no_inertia, eng, timing, false) } +fn turn_src(p: &Persona, st: &State, raw: &Json, no_inertia: bool, eng: Engine, timing: bool, check: bool) -> R<(Json, State)> { let raw: Vec<(String, Json)> = match raw { Json::Null => vec![], Json::Obj(v) => v.clone(), _ => return Err(perr("inputs", "must be a JSON object")) }; for (n, (k, _)) in raw.iter().enumerate() { if raw[..n].iter().any(|(k2, _)| k2 == k) { return Err(perr(&format!("inputs.{k}"), "duplicate input")); } } + // reward provenance (§2.10), first: with `reward_from`, `src` is read here (not compiled, so not `ignored`) and comes back + // in the stance's inputs; without it `src` is an undeclared input as before + let src = match &p.reward { Some(rw) => { if check { provenance(rw, &raw)?; } get(&raw, "src").cloned() } None => None }; + let raw: Vec<(String, Json)> = if p.reward.is_some() { raw.into_iter().filter(|(k, _)| k != "src").collect() } else { raw }; // a drives block steps the drives before the turn, which runs on the stepped state with the drive aggregates as inputs let given = raw; let (st2, raw) = prep(p, st, &given)?; let st = &st2; let t0 = Instant::now(); let mut calls = 0; @@ -1346,6 +1406,7 @@ pub fn turn(p: &Persona, st: &State, raw: &Json, no_inertia: bool, eng: Engine, ("engine_calls".into(), Json::Num(calls as f64)), ("holds".into(), Json::Num(out.held.len() as f64)), ("load_avg_1m".into(), crate::sys::loadavg().map_or(Json::Null, |l| Json::Num((l * 100.0).round() / 100.0)))])); } let mut j = out.to_json(p, st); if let (Some(d), Some(ds)) = (&p.drives, &st.drives) { stance_drives(&mut j, d, ds, &out.lift, &given); } + if let (Some(s), Json::Obj(v)) = (src, &mut j) { if let Some((_, Json::Obj(ins))) = v.iter_mut().find(|(k, _)| k == "inputs") { ins.push(("src".into(), s)); } } Ok((j, ns)) } @@ -1579,6 +1640,8 @@ pub fn describe(p: &Persona) -> Json { v.push(("drives".into(), Json::Obj(vec![("goals".into(), Json::Arr(dr.goals.iter().map(|g| s(&g.id)).collect())), ("say".into(), per(&|g| Some(s(&g.say)))), ("floor".into(), per(&|g| (g.floor > 0.0).then_some(Json::Num(g.floor)))), ("starve_after".into(), per(&|g| g.starve_after.map(Json::Num))), ("learn_from_surprise".into(), Json::Num(dr.kappa))]))); } + // reward_from (§2.10), when declared: per reward-bearing input its allowed source kinds, or "any" + if let (Json::Obj(v), Some(rw)) = (&mut d, &p.reward) { v.push(("reward_from".into(), Json::Obj(rw.iter().map(|(k, a)| (k.clone(), a.as_ref().map_or(s("any"), |a| jstrs(a)))).collect()))); } d } /// A trait or input the drives block adds when it is lowered (`pursue`, the drive aggregates and goal conditions): not the author's diff --git a/probbit-cli/tests/monitor.rs b/probbit-cli/tests/monitor.rs index cb3b99c..cc217d8 100644 --- a/probbit-cli/tests/monitor.rs +++ b/probbit-cli/tests/monitor.rs @@ -11,7 +11,7 @@ use std::time::{Duration, Instant}; const WEEK: &str = "tests/fixtures/monitor/tutor-week.strand"; fn probbit(args: &[&str]) -> (i32, String, String) { - let o = Command::new(env!("CARGO_BIN_EXE_probbit")).args(args).env_remove("NO_COLOR").env_remove("PROBBIT_THEME").stdin(Stdio::null()).output().unwrap(); + let o = Command::new(env!("CARGO_BIN_EXE_probbit")).args(args).env_remove("NO_COLOR").env_remove("PROBBIT_THEME").env_remove("TERM").stdin(Stdio::null()).output().unwrap(); (o.status.code().unwrap_or(-1), String::from_utf8_lossy(&o.stdout).into_owned(), String::from_utf8_lossy(&o.stderr).into_owned()) } fn tmp(tag: &str) -> String { std::env::temp_dir().join(format!("probbit-monitor-{tag}-{}.strand", std::process::id())).to_str().unwrap().to_string() } diff --git a/probbit-cli/tests/persona.rs b/probbit-cli/tests/persona.rs index 5834667..42f2e39 100644 --- a/probbit-cli/tests/persona.rs +++ b/probbit-cli/tests/persona.rs @@ -642,6 +642,74 @@ fn live_watch_follows_an_events_file() { for x in [&s1, &s2] { let _ = std::fs::remove_file(x); } } +// ---------------- the safety kit (docs/persona.md §2.10, §5.7) +/// 300 events of the drives fixture with `src` on each (self, a person, a sensor, the clock), rewards on many +fn src_events() -> Vec { + (0..300).map(|t| { let mut e = vec![r#""elapsed_hours":0.5"#.to_string()]; + if t % 2 == 0 { e.push(r#""praise":true"#.into()); } if t % 5 == 1 { e.push(r#""criticism":true"#.into()); } if t % 3 == 0 { e.push(r#""goals":{"fun":{"win":1.0}}"#.into()); } + e.push(format!(r#""src":"{}""#, ["self", "human:owner", "env:tests", "clock"][t % 4])); format!("{{{}}}", e.join(",")) }).collect() +} +/// A persona without `reward_from` gives 0.8.0's bytes with `src` in its events (an undeclared input, listed in `ignored`): the +/// trace, the strand and the final state of 300 events of the drives fixture (seed 4) are the ones the 0.8.0 binary writes, and +/// the run leaves no lock file behind. The same persona with `reward_from: [human, env]` refuses the rewards from `self` and the +/// clock (exit 2, one error at `inputs.src` per refused event) and accepts the others. +#[test] +fn src_without_reward_from_gives_0_8_0_bytes() { + let fx = format!("{}/tests/fixtures/persona/drives-adversary.json", env!("CARGO_MANIFEST_DIR")); + let (script, st, sd) = (tmp("src-script.json"), tmp("src-state.json"), tmp("src.strand")); for f in [&st, &sd] { let _ = std::fs::remove_file(f); } + let evs = src_events(); std::fs::write(&script, format!("[{}]", evs.join(","))).unwrap(); + let (c, _, e) = probbit(&["persona", "replay", &fx, "--seed", "4", "--script", &script], ""); assert_eq!(c, 0, "{e}"); + assert_eq!(e.trim(), "replay: 300 turns, trace sha256 ce8f247d1fe761ea651998f48360491030f15ba3c8530b0ed59d824868f71b71, final state sha256:52312b3cb8744360e7613e4db2fc5e3d326cb3c2204af7312e2a4196f9781551"); + let (_, s0, _) = probbit(&["persona", "init", &fx, "--seed", "4"], ""); std::fs::write(&st, &s0).unwrap(); + let (c, out, e) = probbit(&["live", &fx, "--state", &st, "--clock", "fixed", "--strand", &sd], &(evs.join("\n") + "\n")); assert_eq!(c, 0, "{e}"); + assert!(e.contains("strand head sha256:c37fea55abbec95d0c0f1f1703741e20c26cd5d66587132917772435a5e94655, final state sha256:52312b3cb8744360e7613e4db2fc5e3d326cb3c2204af7312e2a4196f9781551"), "{e}"); + assert_eq!(out.lines().filter(|l| l.contains(r#""ignored":["src"]"#)).count(), 300); + assert!(!std::path::Path::new(&format!("{sd}.lock")).exists(), "the lock is released"); + // with reward_from: the rewards from self and the clock are refused, the others accepted + let mut d = parse(&read(&fx)); if let Json::Obj(kv) = &mut d { kv.push(("reward_from".into(), parse(r#"["human","env"]"#))); } + let guarded = tmp("src-guarded.json"); std::fs::write(&guarded, jw(&d)).unwrap(); let gs = tmp("src-guarded-state.json"); + let (_, s0, _) = probbit(&["persona", "init", &guarded, "--seed", "4"], ""); std::fs::write(&gs, &s0).unwrap(); + let (c, out, _) = probbit(&["live", &guarded, "--state", &gs, "--clock", "fixed"], &(evs.join("\n") + "\n")); assert_eq!(c, 2); + let refused = out.lines().filter(|l| l.contains(r#""path":"events["#) && l.contains(r#"].inputs.src""#)).count(); + let rewarded = |t: usize| t % 2 == 0 || t % 5 == 1 || t % 3 == 0; + assert_eq!(refused, (0..300).filter(|t| rewarded(*t) && (t % 4 == 0 || t % 4 == 3)).count()); + let (c, out, _) = probbit(&["persona", "turn", &guarded, "--state", &gs, "--inputs", r#"{"praise":true,"src":"self"}"#], ""); assert_eq!(c, 2); assert!(out.contains(r#""path":"inputs.src""#), "{out}"); + for f in [&script, &st, &sd, &guarded, &gs] { let _ = std::fs::remove_file(f); } +} + +/// The safety kit through the CLI: `--checkpoint-every` writes checkpoint lines that `verify` and `verify --from-checkpoint` +/// check; `live control` pauses (every event refused, exit 4, nothing written), resumes and retires (final); the individual cannot +/// control itself; a held lock refuses a second writer and a control line (exit 4) and changes nothing; no lock is left behind +#[test] +fn live_safety_kit_exit_codes() { + let (st, sd) = (tmp("kit-state.json"), tmp("kit.strand")); let lock = format!("{sd}.lock"); for f in [&st, &sd, &lock] { let _ = std::fs::remove_file(f); } + let tutor = ex("tutor.yaml"); + let (_, s0, _) = probbit(&["persona", "init", &tutor, "--seed", "3"], ""); std::fs::write(&st, &s0).unwrap(); + let evs = |n: usize| (0..n).map(|t| format!(r#"{{"praise":{},"elapsed_hours":1}}"#, t % 2 == 0)).collect::>().join("\n") + "\n"; + let (c, _, e) = probbit(&["live", &tutor, "--state", &st, "--clock", "fixed", "--strand", &sd, "--checkpoint-every", "2"], &evs(5)); assert_eq!(c, 0, "{e}"); + assert_eq!(read(&sd).lines().filter(|l| l.starts_with(r#"{"checkpoint":"#)).count(), 2); + let (c, v, _) = probbit(&["live", "verify", &sd], ""); assert_eq!(c, 0, "{v}"); assert!(v.contains(r#""checkpoints":2"#), "{v}"); + let (c, v, _) = probbit(&["live", "verify", &sd, "--from-checkpoint"], ""); assert_eq!(c, 0, "{v}"); assert!(v.contains(r#""from_checkpoint":4"#), "{v}"); + let (c, o, e) = probbit(&["live", "control", &sd, "pause", "--by", "human:owner", "--reason", "a check", "--at", "2026-10-08T10:00:00Z"], ""); assert_eq!(c, 0, "{e}"); + assert!(o.contains(r#""status":"paused""#), "{o}"); + let before = read(&sd); + let (c, o, _) = probbit(&["live", &tutor, "--state", &st, "--clock", "fixed", "--strand", &sd], &evs(3)); assert_eq!(c, 4); + assert_eq!(o.lines().filter(|l| l.contains(r#""code":"paused""#)).count(), 3); assert_eq!(read(&sd), before, "nothing written while paused"); + let (c, _, _) = probbit(&["live", "control", &sd, "resume", "--by", "self", "--reason", "x"], ""); assert_eq!(c, 2, "the individual cannot resume itself"); + let (c, _, e) = probbit(&["live", "control", &sd, "resume", "--by", "human:owner", "--reason", "checked"], ""); assert_eq!(c, 0, "{e}"); + let (c, _, e) = probbit(&["live", &tutor, "--state", &st, "--clock", "fixed", "--strand", &sd], &evs(2)); assert_eq!(c, 0, "{e}"); + // a lock held by a running process (this test's): a second writer and a control line exit 4 and change nothing + std::fs::write(&lock, format!(r#"{{"pid":{},"since":"2026-10-08T10:00:00Z","t":1}}"#, std::process::id())).unwrap(); let before = read(&sd); + let (c, o, _) = probbit(&["live", &tutor, "--state", &st, "--clock", "fixed", "--strand", &sd], &evs(1)); assert_eq!(c, 4); assert!(o.contains(r#""code":"locked""#), "{o}"); + let (c, _, _) = probbit(&["live", "control", &sd, "pause", "--by", "human:owner", "--reason", "x"], ""); assert_eq!(c, 4, "control takes the lock too"); + assert_eq!(read(&sd), before); std::fs::remove_file(&lock).unwrap(); + let (c, _, _) = probbit(&["live", "control", &sd, "retire", "--by", "human:owner", "--reason", "the end"], ""); assert_eq!(c, 0); + let (c, _, _) = probbit(&["live", "control", &sd, "resume", "--by", "human:owner", "--reason", "again"], ""); assert_eq!(c, 4, "retire is final"); + let (c, v, _) = probbit(&["live", "verify", &sd], ""); assert_eq!(c, 0, "{v}"); assert!(v.contains(r#""status":"retired""#) && v.contains(r#""controls":3"#), "{v}"); + assert!(!std::path::Path::new(&lock).exists(), "no lock left behind"); + for f in [&st, &sd] { let _ = std::fs::remove_file(f); } +} + /// The `drives` block (0.8.0) is additive: a persona without it gives 0.7.0's documents byte for byte, with `goals` in the inputs /// too (an ignored input there, as it was in 0.7.0). The digests below are 0.7.0's own `persona replay` output on the same scripts /// (tests/fixtures/persona/drives-noblock.json: 40 random turns per example persona, with goal signals and idle hours).