diff --git a/docs/tui.md b/docs/tui.md index 1d338229..06ea66f0 100644 --- a/docs/tui.md +++ b/docs/tui.md @@ -44,7 +44,8 @@ needing cmake and a C++17 compiler). This is also why the TUI is not part of | `/` + a command | a slash command — see below; `Tab` completes the name, `/help` lists them | | `Ctrl-J` | a new line in the compose box, which grows a row per line, typed or soft-wrapped at its width (up to `:max-lines`, 8, then it scrolls); `Enter` sends | | `Ctrl-P` / `Ctrl-N` | walk back and forth through what you sent | -| `PgUp` / `PgDn`, mouse wheel | scroll the conversation; it stops following the bottom | +| `PgUp` / `PgDn` | page the conversation; it stops following the bottom | +| mouse wheel | scroll whichever pane is under the pointer — the conversation, the activity log | | `End`, or `↓` while scrolled up | back to following the bottom | | `Ctrl-O` | open the newest folded result or thinking; again to shut it | | click a fold | open a thinking block or the rest of a result | @@ -193,7 +194,11 @@ One timeline per branch of who said what, in the order it happened (samizdat.tui.timeline), each role in its own voice and colour: - `` — the problem, and every steer a person sent; -- `` — what the model said; its thinking folds under `◇ thinking`; +- `` — what the model said, drawn as markdown (headings, bullets, + code, quotes), with the call syntax it wrote left to the chamber below and + its thinking — the provider's reasoning and any `` block — folded + under `◇ thinking`. While a call streams, the reply appears here as the + model writes it, and gives way to the finished turn when it lands; each tool call is a **chamber**, headed by the tool and the argument it is known by, showing the first `:result-lines` of what came back with the rest a click (or `Ctrl-O`) away, diffs coloured, failures red; @@ -205,11 +210,11 @@ Which journal notes appear, as whom, and where in each note's data its words are, is `tui.edn :conversation :notes`; the branch fetch asks the server for exactly those kinds (`?notes=…`). -It **follows the bottom**: the newest entry holds the frame's focus, so the -pane scrolls as the run speaks. Scrolling up anchors it on an entry, which -stays put however much arrives below; `End` follows again. It is bounded — -the newest `:turns` of them, 60 by default, since every entry is rebuilt on -every frame. +It **follows the bottom**, and scrolls by rows (ftxui-jolt's `:scroll`): the +wheel over it or `PgUp` moves it, and it stays where it was left however much +arrives below; `End`, or sending something, follows again. The activity log +scrolls the same way. It is bounded — the newest `:turns` of them, 60 by +default, since every entry is rebuilt on every frame. ### The footer @@ -370,7 +375,10 @@ is being used, so every failure here is a rendering: The run on screen is **pushed**. The TUI follows `GET /v1/runs/:id/events`, a server-sent event stream (samizdat.api.stream): every journal event after the cursor, then each one as it lands, plus the manifest steps and approval -changes that are never journalled. An event says what changed — a turn on the +changes that are never journalled, and the reply a branch is writing while +its model call streams (`delta` events, published every gates.edn +`:delta-publish-ms` by samizdat.agent.infer, each piece with the offset it +starts at). An event says what changed — a turn on the branch being read, a question for a person, the run ending — and only that is fetched, a burst of events coalesced into one fetch of each thing. The footer says `live` while the stream is up. A dropped stream reconnects with diff --git a/resources/gates.edn b/resources/gates.edn index b91aeb28..fe883924 100644 --- a/resources/gates.edn +++ b/resources/gates.edn @@ -2060,6 +2060,16 @@ Newest first, so what the cap drops is the oldest change — the one least likely to be what someone is looking at."} + :delta-publish-ms + {:value 100 + :kind :cost-ceiling :capability-tunable? false + :provenance ["karamazov-ta8w"] + :doc "How often a reply being written is published to the event bus, in + ms, while a run's model call streams (samizdat.llm.stream). The + pieces between publishes are sent together. Lower is a livelier + conversation pane and more events; the bus keeps a window of 256 per + watcher, and a watcher that falls behind loses the oldest."} + :event-stream {:value {:poll-ms 100 :heartbeat-ms 15000 :page 200} :kind :cost-ceiling :capability-tunable? false diff --git a/resources/tui.edn b/resources/tui.edn index 27b2117a..acce1a8e 100644 --- a/resources/tui.edn +++ b/resources/tui.edn @@ -83,6 +83,12 @@ :system {:color "#6a8c78"} :thinking {:color "#80968c"} + ;; a reply's markdown + :md-h1 {:color "#9eeeac" :bold true} + :md-h2 {:color "#60dc8c" :bold true} + :md-code {:color "#769e84"} + :md-quote {:color "#567460"} + ;; what the agent did :tool {:color "#6cbc96"} :tool-box {:border :rounded :border-color "#486252"} @@ -303,10 +309,12 @@ ;; The avatar and the compose box. No height: the strip is as tall as the ;; box, which grows a row per line — typed (Ctrl+J) or wrapped — up to - ;; :max-lines. + ;; :max-lines. The box runs from under the conversation to the right edge: + ;; its minimum is the middle and right columns' together (120 + 34), and + ;; it grows alongside the avatar, so the two split the spare width the way + ;; the left and right columns above do and the edges line up. [:hbox [:widget/avatar {:width [:>= 32] :flex :grow}] - [:widget/input {:width 120 :flex :shrink :boxed true :max-lines 8}] - [:filler {:width [:>= 34] :flex :grow}]] + [:widget/input {:width [:>= 154] :flex :grow :boxed true :max-lines 8}]] [:widget/status {}]]} diff --git a/src/samizdat/agent/beam.clj b/src/samizdat/agent/beam.clj index 8f7ab396..080f34f1 100644 --- a/src/samizdat/agent/beam.clj +++ b/src/samizdat/agent/beam.clj @@ -1050,7 +1050,10 @@ ;; a moment; the journal poller handles that, and it is the honest ;; picture — the branches genuinely do not exist yet. run-id (runs/start-run! conn {:problem problem - :provider (:provider llm-config) + ;; The alias, so a resume resolves the + ;; same declaration (api.control/resume!). + :provider (or (:provider-name llm-config) + (:provider llm-config)) :model (:model llm-config) :max-turns max-turns :beam-width requested-width diff --git a/src/samizdat/agent/infer.clj b/src/samizdat/agent/infer.clj index 6dce03f2..5372b2c4 100644 --- a/src/samizdat/agent/infer.clj +++ b/src/samizdat/agent/infer.clj @@ -55,6 +55,7 @@ cell's job, in resources, where the supervisor can rewrite it." (:require [clojure.tools.logging :as log] [samizdat.agent.gates :as gates] + [samizdat.events :as events] [samizdat.cancel :as cancel] [samizdat.agent.tools.base :as tools] [samizdat.config :as config] @@ -319,6 +320,34 @@ (let [{:keys [turn forced]} (gates/threshold :local-reasoning-budget)] (if force-tool forced turn)))) +(defn- delta-publisher + "An `on-delta` for one call on `branch-id` of `run-id`: the reply's text and + reasoning gathered and published to the bus at most every gates.edn + :delta-publish-ms, each piece with the offset it starts at, and a `flush!` + for what is still gathered when the call returns. {:on-delta :flush!}." + [run-id branch-id] + (let [every (gates/threshold :delta-publish-ms) + st (atom {:text "" :reasoning "" :sent-text 0 :sent-reasoning 0 :at 0}) + publish! (fn [{:keys [text reasoning sent-text sent-reasoning]}] + (let [t (subs text sent-text) r (subs reasoning sent-reasoning)] + (when (or (seq t) (seq r)) + (events/publish! (cond-> {:kind :delta :run-id run-id :branch-id branch-id} + (seq t) (assoc :text t :text-at sent-text) + (seq r) (assoc :reasoning r :reasoning-at sent-reasoning)))))) + mark (fn [s now] (-> s + (assoc :sent-text (count (:text s)) + :sent-reasoning (count (:reasoning s)) :at now) + (update :marks (fnil inc 0))))] + {:on-delta (fn [{:keys [text reasoning]}] + (let [now (System/currentTimeMillis) + [before after] (swap-vals! st #(cond-> (-> % (update :text str text) + (update :reasoning str reasoning)) + (>= (- now (:at %)) every) (mark now)))] + (when (not= (:marks before) (:marks after)) + (publish! (assoc after :sent-text (:sent-text before) + :sent-reasoning (:sent-reasoning before)))))) + :flush! (fn [] (let [[before _] (swap-vals! st mark 0)] (publish! before)))})) + (defn complete-fn "ctx -> (fn [tape] -> {:ok true :response r} | {:ok false :error s}). @@ -351,7 +380,12 @@ (force-mechanism ctx tape) reasoning-budget (reasoning-budget-for ctx tape)] (loop [attempt 1] - (let [base (or (:max-tokens (:llm-config ctx)) + (let [;; Somebody may be watching the run: its calls stream, and what + ;; the model writes goes onto the bus as it is written. Not a + ;; probe's (journal? false): a bounce is not a turn anyone sees. + watch (when (and journal? (:run-id ctx) id) + (delta-publisher (:run-id ctx) (str id))) + base (or (:max-tokens (:llm-config ctx)) ;; No configured cap: the FIRST attempt keeps the ;; provider's default, but a retry exists to buy room — ;; doubling nothing was a same-budget repeat (blt.38). @@ -377,7 +411,8 @@ ;; The stable conversation key an endpoint ;; pins its prefix cache to. Only the local ;; adapter emits it; see LR-5. - id (assoc :cache-key (str id)))) + id (assoc :cache-key (str id)) + watch (assoc :on-delta (:on-delta watch)))) :wire wire)} (catch Throwable e ;; The reason travels with the failure. `provider-error-step` @@ -386,7 +421,8 @@ ;; versus wait and retry. Without this the loop knows only ;; `the call failed` and every provider problem looks alike. {:ok false :error (ex-message e) - :reason (or (:reason (ex-data e)) :call-failed)}))] + :reason (or (:reason (ex-data e)) :call-failed)}) + (finally (some-> watch :flush! (apply []))))] (if (and (:ok r) (< attempt max-call-attempts) ;; The prefill the adapter ACTUALLY sent (nil where it was diff --git a/src/samizdat/agent/resume.clj b/src/samizdat/agent/resume.clj index cff45303..3a49f259 100644 --- a/src/samizdat/agent/resume.clj +++ b/src/samizdat/agent/resume.clj @@ -400,7 +400,9 @@ :title (:title held)}) held (task-tool/task-statement held) (seq in-flight) (update :messages into in-flight)))) - (runs/branches conn run-id)) + ;; Not the supervisor stream's branch: the stream + ;; carries it and re-opens it on its next pass. + (remove #(= "supervisor" (:role %)) (runs/branches conn run-id))) ;; The anchor: rounds completed are the max turn in the journal, so ;; the loop continues one past it. max-turns is the ORIGINAL budget. start-turn (inc (reduce max 0 (map :turn turn-rows)))] diff --git a/src/samizdat/api/control.clj b/src/samizdat/api/control.clj index 26419003..dcb2b341 100644 --- a/src/samizdat/api/control.clj +++ b/src/samizdat/api/control.clj @@ -223,15 +223,37 @@ `body` may carry max_turns: an explicit budget extension that reopens branches closed as exhausted. Omitted, the original budget stands." [{:keys [conn config]} run-id body] - (if-not (resume/resumable? conn run-id) - {:status 409 :body {:error {:message (str "run " run-id " is not resumable") - :run_id run-id}}} + (let [refuse (fn [why] {:status 409 :body {:error {:message (str "run " run-id " " why) + :run_id run-id}}}) + run (runs/get-run conn run-id) + ;; The model the run was ON, from its row: the provider it recorded + ;; (an alias config.edn declares, or a built-in) and the model. + recorded (when (not-empty (str (:provider run))) + (try (config/provider-llm config (:provider run) + (if (not-empty (str (:model run))) + {:model (:model run)} + {})) + (catch Exception e e)))] + (cond + (not (resume/resumable? conn run-id)) (refuse "is not resumable") + (instance? Exception recorded) (refuse (str "ran on " (:provider run) ": " + (ex-message recorded))) + :else ;; A resume may name an arm too — a run that crashed on one model can be ;; picked up on another, and saying nothing keeps the original. - (let [llm-config (run-llm-config config (:llm config) body) + (let [llm-config (run-llm-config config (or recorded (:llm config)) body) adapter (registry/adapter-for (:provider llm-config)) abort (atom false) - max-turns (or (:max_turns body) (:max-turns body))] + max-turns (or (:max_turns body) (:max-turns body)) + ;; A run already driven by this process is not resumed: a second + ;; driver over the same branches wrote duplicate turns, every call + ;; on whatever model the resume resolved. Claimed HERE, atomically, + ;; not by the spawned thread, or two resumes a double click apart + ;; would both get in before either thread registered. + [before _] (swap-vals! active #(if (contains? % run-id) % (assoc % run-id {:abort abort})))] + (if (contains? before run-id) + (refuse "is still running") + (do (let [cancel* (atom nil) started (cancel/start! (cancel/spawn @@ -261,7 +283,7 @@ ;; max_turns extension was reported as the old budget more often than ;; not. {:body {:run_id run-id :status "resuming" - :max_turns (or max-turns (:max_turns (runs/get-run conn run-id)))}}))) + :max_turns (or max-turns (:max_turns (runs/get-run conn run-id)))}})))))) (defn- grant-pattern "The pattern from a grant payload. Accepts a map (what body-json yields), a bare string, or nil. Blank is not a pattern — an unset form posts empty diff --git a/src/samizdat/api/sse.clj b/src/samizdat/api/sse.clj index da5b7c5a..39bb5fed 100644 --- a/src/samizdat/api/sse.clj +++ b/src/samizdat/api/sse.clj @@ -38,10 +38,17 @@ ;; --- the parser -------------------------------------------------------------- (defn reader - "A fresh parser: feed it the bytes of one connection with `feed`." - [] - {:phase :head :buf [] :status nil :chunked? false - :need nil :skip 0 :line [] :event {} :done? false}) + "A fresh parser: feed it the bytes of one connection with `feed`. + + `:keep-body? true` keeps a response that is not an event stream — a status + other than 200, or a 200 with some other content type — as its raw body + (`:raw`, `body-text`) rather than giving it up: a provider's error, or a + provider that answered a streamed request whole." + ([] (reader nil)) + ([{:keys [keep-body?]}] + {:phase :head :buf [] :status nil :headers {} :chunked? false + :need nil :skip 0 :line [] :event {} :done? false + :keep-body? (boolean keep-body?) :raw? false :raw []})) (defn- utf8 [bs] (String. (byte-array bs) "UTF-8")) @@ -94,6 +101,11 @@ [(update st :line conj b) out])) [st out] bs)) +(defn- sink + "Body bytes to where they go: the raw body, or the event lines." + [st out bs] + (if (:raw? st) [(update st :raw into bs) out] (payload st out bs))) + (defn- chunked "De-chunk `bs` into the event parser." [st out bs] @@ -115,13 +127,26 @@ :else (let [k (min (:need st) (count bs)) - [st out] (payload st out (take k bs)) + [st out] (sink st out (take k bs)) left (- (:need st) k)] (recur (if (zero? left) (assoc st :need nil :skip 2) (assoc st :need left)) out (seq (drop k bs))))))) (defn- body [st out bs] - (if (:chunked? st) (chunked st out bs) (payload st out bs))) + (if (:chunked? st) (chunked st out bs) (sink st out bs))) + +(defn- head-fields + "The response head's header fields, names lower-cased." + [head] + (into {} (for [line (rest (str/split-lines head)) + :let [i (str/index-of line ":")] + :when i] + [(str/lower-case (str/trim (subs line 0 i))) (str/trim (subs line (inc i)))]))) + +(defn body-text + "A kept body (`:keep-body?`) as text." + [st] + (utf8 (:raw st))) (defn feed "Feed `bytes` to parser state `st`. Returns {:state :events}: the events @@ -135,9 +160,15 @@ (if-let [end (head-end buf)] (let [head (utf8 (subvec buf 0 end)) status (some-> (re-find #"^HTTP/\d\.\d (\d{3})" head) second parse-long) - chunked? (boolean (re-find #"(?i)\r\ntransfer-encoding:\s*chunked" head)) - st (assoc st :phase :body :buf [] :status status :chunked? chunked?)] - (if (= 200 status) + fields (head-fields head) + chunked? (boolean (some-> (get fields "transfer-encoding") str/lower-case + (str/includes? "chunked"))) + events? (boolean (some-> (get fields "content-type") str/lower-case + (str/includes? "text/event-stream"))) + raw? (and (:keep-body? st) (or (not= 200 status) (not events?))) + st (assoc st :phase :body :buf [] :status status :headers fields + :chunked? chunked? :raw? raw?)] + (if (or raw? (= 200 status)) (let [[st out] (body st [] (subvec buf end))] {:state st :events out}) {:state (assoc st :done? true) :events []})) @@ -154,9 +185,13 @@ ;; --- following a stream -------------------------------------------------------- -(defn- parse-url [url] - (let [[_ host port path] (re-find #"^https?://([^:/]+)(?::(\d+))?(/.*)?$" (str url))] - {:host host :port (or (some-> port parse-long) 80) :path (or path "/")})) +(defn parse-url + "`url` as {:tls? :host :port :path}, the port defaulting by scheme." + [url] + (let [[_ scheme host port path] (re-find #"^(https?)://([^:/]+)(?::(\d+))?(/.*)?$" (str url)) + tls? (= "https" scheme)] + {:tls? tls? :host host :port (or (some-> port parse-long) (if tls? 443 80)) + :path (or path "/")})) (defn- request-text [{:keys [host port path]} last-id] (str "GET " path " HTTP/1.1\r\n" diff --git a/src/samizdat/config.clj b/src/samizdat/config.clj index 428c9eaa..575572bc 100644 --- a/src/samizdat/config.clj +++ b/src/samizdat/config.clj @@ -257,7 +257,7 @@ ;; URL continues a flagged assistant prefix; `thinking {type ;; disabled}` reliably yields no reasoning; a native tool_choice ;; is honoured, but only with thinking off (the adapter knows). - :features #{:prefill :native-tool-choice :thinking-toggle :reasoning-effort} + :features #{:prefill :native-tool-choice :thinking-toggle :reasoning-effort :stream} :key-env "DEEPSEEK_API_KEY" ;; deepseek-v4-flash is the development and test model: cheap ;; enough to run the beam repeatedly. deepseek-v4-pro is the @@ -275,7 +275,7 @@ ;; Measured 2026-09-06: ignores a trailing assistant prefix ;; entirely, cannot be told not to think (effort low is the ;; least), honours tool_choice {type function} despite its docs. - :features #{:native-tool-choice :reasoning-effort} + :features #{:native-tool-choice :reasoning-effort :stream} :key-env "ZHIPU_API_KEY" :model "glm-5.3" ;; GLM benefits from a low temperature on coding tasks (dirge @@ -283,7 +283,7 @@ :temperature 0.2} :openai {:context-window 128000 :base-url "https://api.openai.com/v1" - :features #{:native-tool-choice :reasoning-effort} + :features #{:native-tool-choice :reasoning-effort :stream} :key-env "OPENAI_API_KEY" :model "gpt-4o"} ;; A local llama-server / vLLM / LM Studio OpenAI-compatible endpoint. @@ -297,7 +297,7 @@ ;; A bare OpenAI-compatible endpoint until the startup probe ;; says which server it is; `llama-cpp-features` is what a ;; llama.cpp answer adds (apply-discovery). - :features #{:native-tool-choice :reasoning-effort}} + :features #{:native-tool-choice :reasoning-effort :stream}} ;; Ollama's NATIVE api, so no /v1 suffix. See llm/adapter/ollama.clj for ;; why the native surface rather than Ollama's OpenAI-compatible one. :ollama {:context-window 32768 @@ -327,6 +327,8 @@ ;; thinking {type disabled}, llama.cpp's ;; chat_template_kwargs {enable_thinking false}) ;; :reasoning-effort a top-level reasoning_effort is honoured +;; :stream `stream: true` is answered as server-sent chunks, so +;; a reply can be watched as it is written (def llama-cpp-features "What a llama.cpp server adds once /props has identified it. Measured on diff --git a/src/samizdat/llm/client.clj b/src/samizdat/llm/client.clj index 5c63ce2a..7f39439c 100644 --- a/src/samizdat/llm/client.clj +++ b/src/samizdat/llm/client.clj @@ -49,6 +49,7 @@ [samizdat.llm.adapter :as adapter] [samizdat.llm.message :as message] [samizdat.llm.ratelimit :as ratelimit] + [samizdat.llm.stream :as stream] [samizdat.util :as util] [samizdat.session :as session])) @@ -244,16 +245,21 @@ ;; storing a doubled fence on every steered GLM turn. use-prefill? (boolean (and (:prefill request) (adapter/prefill-support? adapter config))) - body (adapter/chat-body adapter config request) + streamed? (boolean (and (:on-delta request) (contains? (:features config) :stream))) + body (cond-> (adapter/chat-body adapter config request) + ;; Usage rides the last chunk only when asked for. + streamed? (assoc :stream true :stream_options {:include_usage true})) payload (json/write-str body) started (System/currentTimeMillis) - resp (http/post url {:headers (merge (request-headers adapter config) - {"Content-Type" "application/json"}) - :body payload - :socket-timeout (:timeout-ms config default-timeout-ms) - :conn-timeout (:conn-timeout-ms config - default-conn-timeout-ms) - :throw-exceptions false}) + opts {:headers (merge (request-headers adapter config) + {"Content-Type" "application/json"}) + :body payload + :socket-timeout (:timeout-ms config default-timeout-ms) + :conn-timeout (:conn-timeout-ms config default-conn-timeout-ms)} + resp (if streamed? + (stream/post url (assoc opts :max-response-ms (:max-response-ms config)) + (:on-delta request)) + (http/post url (assoc opts :throw-exceptions false))) elapsed (- (System/currentTimeMillis) started) status (:status resp) decoded (decode (:body resp))] @@ -371,7 +377,8 @@ stuck provider costs a known amount rather than the run." ([adapter config messages] (chat adapter config messages nil)) ([adapter config messages {:keys [max-tokens temperature max-retries prefill force-tool - cache-key reasoning-effort grammar reasoning-budget]}] + cache-key reasoning-effort grammar reasoning-budget + on-delta]}] (let [request {:messages (message/prepare messages) :max-tokens (or max-tokens (:max-tokens config)) :temperature (or temperature (:temperature config)) @@ -402,7 +409,11 @@ :grammar grammar ;; A per-call thinking cap for a llama.cpp endpoint ;; (karamazov-w7n4); every other adapter ignores it. - :reasoning-budget reasoning-budget} + :reasoning-budget reasoning-budget + ;; Somebody watching the reply as it is written: with an + ;; endpoint that can stream (:stream), each delta goes to + ;; it as it arrives (samizdat.llm.stream). + :on-delta on-delta} ;; The read timeout is sized to the budget being asked for: a big ;; max-tokens legitimately takes longer than a small one, and a fixed ;; bound cut off long generations and re-billed them (see diff --git a/src/samizdat/llm/fence.clj b/src/samizdat/llm/fence.clj index 0ccfdd4b..f790eeaa 100644 --- a/src/samizdat/llm/fence.clj +++ b/src/samizdat/llm/fence.clj @@ -650,6 +650,29 @@ (str/replace think-re "") (str/replace open-think-re ""))) +;; --- what a reader sees ----------------------------------------------------- + +(def ^:private call-markup-res + "Every call syntax a model writes into its text, each to its closer or to + the end of the reply (a reply cut off mid-call), then the stray closers + some models repeat after a call." + [#"(?s)```tool-call\s*\r?\n.*?(?:```|\z)" + #"(?s).*?(?:|\z)" + #"(?s).*?(?:|\z)" + #"(?s)]*>.*?(?:|\z)" + #"(?s)<\|DSML\|[^>]*>.*?(?:]*>|\z)" + #"]*>"]) + +(defn prose + "The reply as a reader should see it: reasoning blocks and call markup + removed — the thinking fold holds the one and the tool's own chamber shows + the other — and the blank runs the cuts leave collapsed. Display only; the + transcript keeps the reply as it came." + [s] + (-> (reduce #(str/replace %1 %2 "") (strip-think s) call-markup-res) + (str/replace #"\n[ \t]*(?:\n[ \t]*){2,}" "\n\n") + str/trim)) + (defn- parse-tool-call* [response] (let [fenced (extract-fences response) ;; A response that ends in a well-formed call but omits the fence is diff --git a/src/samizdat/llm/stream.clj b/src/samizdat/llm/stream.clj new file mode 100644 index 00000000..de811206 --- /dev/null +++ b/src/samizdat/llm/stream.clj @@ -0,0 +1,173 @@ +;; samizdat - a claim-first verification harness +;; Copyright (C) 2026 Dmitri Sotnikov +;; +;; This program is free software: you can redistribute it and/or modify +;; it under the terms of the GNU General Public License as published by +;; the Free Software Foundation, either version 3 of the License, or +;; (at your option) any later version. +;; +;; This program is distributed in the hope that it will be useful, +;; but WITHOUT ANY WARRANTY; without even the implied warranty of +;; MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +;; GNU General Public License for more details. +;; +;; You should have received a copy of the GNU General Public License +;; along with this program. If not, see . +;; +;; SPDX-License-Identifier: GPL-3.0-or-later + +(ns samizdat.llm.stream + "A chat completion STREAMED from an OpenAI-compatible endpoint: each delta + handed on as it arrives, and the chunks folded back into the one completion + body the adapters already parse — so everything after the transport + (status handling, parsing, retries) is the same code whether a call + streamed or not. + + jolt.http-client reads a response whole, so this speaks HTTP/1.1 over the + socket itself, as samizdat.api.sse does for the harness's own event + stream, and reads the body with that namespace's parser: plain TCP through + jolt.http.net, TLS through jolt.http.tls." + (:require ;; the java.time.* host shim, before data.json — see samizdat.store.journal + [jolt.time] + [clojure.data.json :as json] + [clojure.string :as str] + [jolt.http.net :as net] + [jolt.http.tls :as tls] + [samizdat.api.sse :as sse])) + +;; --- the chunks -------------------------------------------------------------- + +(defn- merge-call + "One streamed tool-call fragment into the call at its index: the id, type + and name as they first come, the arguments concatenated." + [call frag] + (let [f (:function frag)] + (cond-> (or call {}) + (:id frag) (assoc :id (:id frag)) + (:type frag) (assoc :type (:type frag)) + (:name f) (update-in [:function :name] str (:name f)) + (:arguments f) (update-in [:function :arguments] str (:arguments f))))) + +(defn accumulate + "Fold one chunk into `acc`. Every string a delta carries is appended to the + same field of the message — content, and whichever field the endpoint puts + its reasoning in — so no provider's field needs naming here." + [acc chunk] + (let [{:keys [delta finish_reason]} (first (:choices chunk))] + (cond-> (merge acc (select-keys chunk [:id :model :created])) + (:usage chunk) (assoc :usage (:usage chunk)) + finish_reason (assoc :finish-reason finish_reason) + delta + (update :message + (fn [m] + (reduce-kv (fn [m k v] + (cond + (= :tool_calls k) + (reduce #(update-in %1 [:calls (:index %2 0)] merge-call %2) m v) + (string? v) (update m k str v) + (some? v) (assoc m k v) + :else m)) + (or m {}) delta)))))) + +(defn completion + "`acc` as the completion a non-streamed call returns." + [{:keys [message usage finish-reason] :as acc}] + (let [calls (:calls message)] + (cond-> {:id (:id acc) :object "chat.completion" :model (:model acc) + :choices [{:index 0 + :message (cond-> (-> (dissoc message :calls) + (update :role #(or % "assistant"))) + (seq calls) (assoc :tool_calls (mapv val (sort-by key calls)))) + :finish_reason finish-reason}]} + usage (assoc :usage usage)))) + +(defn delta + "What a chunk adds that a reader sees: {:text} and/or {:reasoning}, or nil." + [chunk] + (let [d (:delta (first (:choices chunk))) + text (:content d) + reasoning (or (:reasoning_content d) (:reasoning d))] + (not-empty (cond-> {} + (not-empty (str text)) (assoc :text text) + (not-empty (str reasoning)) (assoc :reasoning reasoning))))) + +;; --- the socket -------------------------------------------------------------- + +(defn- open + "A connection as {:send! :recv! :close!}. `recv!` answers a chunk of bytes, + nil at end of stream, or throws a read timeout." + [{:keys [tls? host port]} read-ms conn-ms] + (if tls? + (let [st (tls/tls-connect host port false read-ms conn-ms)] + {:send! #((jolt.host/ref-get st :write) st %) + :recv! #((jolt.host/ref-get st :read) st nil) + :close! #((jolt.host/ref-get st :close))}) + (let [fd (net/connect host port conn-ms)] + (net/set-read-timeout! fd read-ms) + {:send! #(net/send-bytes fd %) + :recv! #(net/recv-bytes fd) + :close! #(net/close fd)}))) + +(defn- request-bytes [{:keys [host port path]} headers ^String body] + (let [b (.getBytes body "UTF-8")] + (byte-array + (concat + (.getBytes + (str "POST " path " HTTP/1.1\r\n" + "Host: " host ":" port "\r\n" + (apply str (for [[k v] headers] (str k ": " v "\r\n"))) + "Accept: text/event-stream\r\n" + "Content-Length: " (alength b) "\r\n" + "Connection: close\r\n\r\n") + "UTF-8") + b)))) + +(defn post + "POST `body` to `url` asking for a stream, calling `on-delta` with each + delta (see `delta`) as it arrives. Returns {:status :headers :body} like + jolt.http-client's post: the body is the folded completion as JSON, or + whatever the endpoint sent instead — an error, or a whole reply to a + request it did not stream. + + `:socket-timeout` bounds each read; `:max-response-ms` the whole response, + as jolt.http.platform's cap does for an unstreamed call. `on-delta` + throwing does not fail the call: a watcher's trouble is not the model's." + [url {:keys [headers body socket-timeout conn-timeout max-response-ms]} on-delta] + (let [t (sse/parse-url url) + started (System/currentTimeMillis) + deadline (when (and max-response-ms (pos? max-response-ms)) (+ started max-response-ms)) + conn (open t socket-timeout conn-timeout) + hand-on (fn [d] (try (on-delta d) (catch Throwable _ nil)))] + (try + ((:send! conn) (request-bytes t headers (str body))) + (loop [st (sse/reader {:keep-body? true}) acc nil] + (when (and deadline (> (System/currentTimeMillis) deadline)) + (throw (jolt.host/throwable "java.net.SocketTimeoutException" + (str "Response exceeded the total time limit of " + max-response-ms "ms")))) + (let [got ((:recv! conn)) + {:keys [state events]} (if got (sse/feed st got) {:state st :events []}) + [acc done?] (reduce (fn [[acc done?] {:keys [data]}] + (if (= "[DONE]" (str/trim (str data))) + [acc true] + (let [c (try (json/read-str (str data) :key-fn keyword) + (catch Throwable _ nil))] + (if (map? c) + (do (some-> (delta c) hand-on) + [(accumulate acc c) done?]) + [acc done?])))) + [acc false] events)] + (cond + (and (nil? got) (nil? (:status state))) + (throw (jolt.host/throwable "java.net.SocketException" + "Connection closed before a response")) + + (or (nil? got) done? (:done? state)) + {:status (:status state) + :headers (:headers state) + :body (if (:raw? state) + (sse/body-text state) + (json/write-str (completion acc)))} + + :else (recur state acc)))) + (finally ((:close! conn)))))) diff --git a/src/samizdat/workflow.clj b/src/samizdat/workflow.clj index ca6921a1..fa19b321 100644 --- a/src/samizdat/workflow.clj +++ b/src/samizdat/workflow.clj @@ -342,7 +342,10 @@ opening (instr/opening-context root (:block orient) (gates/threshold :instructions)) run-id (runs/start-run! conn {:problem problem - :provider (:provider llm-config) + ;; The alias, so a resume resolves the + ;; same declaration (api.control/resume!). + :provider (or (:provider-name llm-config) + (:provider llm-config)) :model (:model llm-config) :max-turns max-turns :beam-width 1 diff --git a/test-tui/samizdat/tui/mouse_test.clj b/test-tui/samizdat/tui/mouse_test.clj index e7611681..6457d1e1 100644 --- a/test-tui/samizdat/tui/mouse_test.clj +++ b/test-tui/samizdat/tui/mouse_test.clj @@ -247,6 +247,57 @@ (is (= [(keyword label)] @hit) (str "a click at " (pr-str at) " reached " label))))))))) +(deftest the-compose-box-has-the-focus-from-the-start + ;; Typing into a TUI that has just opened must land in the compose box. The + ;; focus used to start on whatever focusable widget the layout reached + ;; first, so what was typed went nowhere until the box had been clicked. + (let [typed (atom nil) + handlers {:start (fn [& _]) :abort (fn []) :resume (fn []) + :input #(reset! typed %) :submit (fn [_]) :toggle (fn [_]) + :decide (fn [& _]) :answer (fn [& _]) :select-run (fn [_]) + :select-branch (fn [_]) :scroll (fn [& _])} + current (requiring-resolve 'samizdat.tui.layout/current) + expand (requiring-resolve 'samizdat.tui.layout/expand) + themed (requiring-resolve 'samizdat.tui.theme/apply-theme) + turns [{:turn 1 :tool_name "shell" :result "a\nb\nc\nd\ne\nf" :category "success"}] + app (fn [] (let [spec (current)] + (themed (:theme spec) + (expand (:layout spec) + (assoc (st/initial "http://x") :on handlers + :run-id "r" :branch-id "B1" + :branch {:turns turns})))))] + (ui/with-screen [s app] + (ui/render-text s 200 45) + (ui/send-char! s "hi") + (is (= "hi" @typed)))) + (testing "and keeps it as the screen fills in around it" + ;; The first frame has no runs and no turns; the polls that follow add + ;; a run menu and a conversation of folds, each of them focusable. + (let [typed (atom nil) + data (atom {}) + handlers {:start (fn [& _]) :abort (fn []) :resume (fn []) + :input #(reset! typed %) :submit (fn [_]) :toggle (fn [_]) + :decide (fn [& _]) :answer (fn [& _]) :select-run (fn [_]) + :select-branch (fn [_]) :scroll (fn [& _])} + current (requiring-resolve 'samizdat.tui.layout/current) + expand (requiring-resolve 'samizdat.tui.layout/expand) + themed (requiring-resolve 'samizdat.tui.theme/apply-theme) + app (fn [] (let [spec (current)] + (themed (:theme spec) + (expand (:layout spec) + (merge (assoc (st/initial "http://x") :on handlers) + @data)))))] + (ui/with-screen [s app] + (ui/render-text s 200 45) + (reset! data {:runs [{:id "r1" :problem "p" :status "completed"}] + :run-id "r1" :branch-id "B1" + :branch {:turns [{:turn 1 :tool_name "shell" :result "a\nb\nc\nd\ne\nf" + :category "success"}]}}) + (ui/refresh! s) + (ui/render-text s 200 45) + (ui/send-char! s "yo") + (is (= "yo" @typed)))))) + (deftest the-whole-shipped-layout-renders-through-the-real-toolkit ;; The layout is userspace hiccup expanded against the registry, and every ;; widget is only ever checked as data. This is the one test that says the @@ -284,22 +335,21 @@ (is (str/includes? frame "samizdat:trunk") "and the footer's project label"))) (deftest the-conversation-follows-the-bottom-in-the-real-toolkit - ;; The data half says the newest entry carries :focus. That is only the - ;; mechanism; the claim is that FTXUI scrolls its frame to it, so a run - ;; that has said more than fits shows what it said LAST. + ;; The pane follows the bottom until it is scrolled (ftxui-jolt's + ;; :scroll): a run that has said more than fits shows what it said LAST. (let [turns (mapv (fn [n] {:turn n :tool_name "read_file" :result (str "result " n)}) (range 1 41)) frame (ui/render-text (w/conversation (assoc (st/initial "b") :branch {:turns turns}) {}) 60 20)] (is (str/includes? frame "result 40") "the newest turn is on screen") (is (not (str/includes? frame "result 1\n")) "and the oldest has scrolled away")) - (testing "scrolled up to an anchor, the frame holds it instead" + (testing "scrolled to the top, the pane holds it there instead" (let [turns (mapv (fn [n] {:turn n :tool_name "read_file" :result (str "result " n)}) (range 1 41)) frame (ui/render-text (w/conversation (assoc (st/initial "b") :branch {:turns turns} - :scroll-anchor "t5/tool") {}) + :scroll {:conversation {:top 0}}) {}) 60 20)] - (is (str/includes? frame "result 5")) + (is (re-find #"result 1\b" frame) "the oldest turn is on screen") (is (not (str/includes? frame "result 40")))))) (deftest a-multi-select-question-ticks-on-a-click diff --git a/test/samizdat/base_test.clj b/test/samizdat/base_test.clj index 843ebb42..3ba7800e 100644 --- a/test/samizdat/base_test.clj +++ b/test/samizdat/base_test.clj @@ -234,8 +234,7 @@ "src/samizdat/api/sse.clj" {:vocabulary {:all "HTTP/1.1 and server-sent-events wire syntax: the status line, the chunked transfer-encoding header, the URL - form. Protocol, not a vocabulary a project chooses."} - :threshold {80 "The default port of an http URL with none. Protocol."}} + form. Protocol, not a vocabulary a project chooses."}} "src/samizdat/api/stream.clj" {:vocabulary {:all "The query syntax of the stream's own cursor parameter diff --git a/test/samizdat/control_test.clj b/test/samizdat/control_test.clj index 822b41e7..7c5b4795 100644 --- a/test/samizdat/control_test.clj +++ b/test/samizdat/control_test.clj @@ -617,26 +617,24 @@ (deftest a-resumed-branch-reopens-on-the-suffix-it-was-opened-on ;; The row records the suffix the cell handed initial-messages (v24), and ;; the rebuild replays it verbatim: a decompose unit keeps its attempt - ;; framing and a supervisor its role text, whatever the manifest's :prompt - ;; says. Before the column every branch came back on the manifest's prompt — - ;; a resumed unit lost its attempt framing, a resumed supervisor opened on - ;; the supervisor system prompt without the supervisor role text - ;; (karamazov-kgvg). A row that recorded NONE gets none, even under a + ;; framing, whatever the manifest's :prompt says. Before the column every + ;; branch came back on the manifest's prompt — a resumed unit lost its + ;; attempt framing (karamazov-kgvg). A row that recorded NONE gets none, even under a ;; manifest that has a :prompt. (with-db [c] (let [rid (runs/start-run! c {:problem "review src/example.clj" :max-turns 10 :beam-width 1})] - (runs/open-branch! c rid {:branch-id "SUP" :role :supervisor - :prompt-suffix "YOU WATCH THE RUN"}) + (runs/open-branch! c rid {:branch-id "U1" :role :implementor + :prompt-suffix "ATTEMPT THIS UNIT"}) (runs/open-branch! c rid {:branch-id "B1"}) (with-redefs [beam/run-rounds (fn [_ branches _] {:branches branches})] (let [bs (:branches (resume/resume! {:conn c :config {:run {:loop "review"}} :llm-adapter :a :llm-config {} :run-id rid})) system (fn [id] (->> bs (filter #(= id (:id %))) first :messages (filter #(= "system" (:role %))) first :content))] - (is (str/includes? (system "SUP") "YOU WATCH THE RUN") + (is (str/includes? (system "U1") "ATTEMPT THIS UNIT") "the suffix it opened on") - (is (not (str/includes? (system "SUP") "CODE REVIEW")) + (is (not (str/includes? (system "U1") "CODE REVIEW")) "and not the manifest's, which it never saw") (is (not (str/includes? (system "B1") "CODE REVIEW")) "a branch recorded as opening on no suffix gets none")))))) @@ -878,3 +876,55 @@ (is (not (system/started?))) (is (nil? (samizdat.userspace/project-root)) "the project it bound is let go") (is (not (samizdat.userspace/files?))))) + +;; --- resume ------------------------------------------------------------------ + +(deftest a-run-running-in-this-process-cannot-be-resumed + ;; A resume of a live run started a SECOND driver over the same branches: + ;; duplicate turn rows, and every call on whatever model the resume picked + ;; (run 3020cbca, endless-flight). + (with-db [c] + (let [rid (runs/start-run! c {:problem "p" :provider :deepseek :model "m"}) + drove (atom 0)] + (swap! api-control/active assoc rid {:abort (atom false)}) + (try + (with-redefs [resume/resume! (fn [_] (swap! drove inc) {:status :completed})] + (let [r (api-control/resume! {:conn c :config {:llm {:provider :local}}} rid {})] + (is (= 409 (:status r))) + (is (str/includes? (str (get-in r [:body :error :message])) "still running"))) + (Thread/sleep 50) + (is (zero? @drove) "no second driver")) + (finally (swap! api-control/active dissoc rid)))))) + +(deftest a-resume-keeps-the-runs-own-provider-and-model + ;; It rebuilt the model from the config's default, so a run started on + ;; deepseek came back on the local endpoint. + (with-db [c] + (let [config {:llm {:provider :local :model "local-model"} + :providers {:flash {:type :deepseek :model "deepseek-v4-flash"}}} + rid (runs/start-run! c {:problem "p" :provider :flash :model "deepseek-v4-flash"}) + seen (promise)] + (with-redefs [resume/resume! (fn [ctx] (deliver seen (:llm-config ctx)) {:status :completed})] + (api-control/resume! {:conn c :config config} rid {}) + (let [llm (deref seen 2000 nil)] + (is (= :deepseek (:provider llm))) + (is (= :flash (:provider-name llm))) + (is (= "deepseek-v4-flash" (:model llm))))) + (testing "a resume that names a model still gets it" + (let [seen (promise)] + (with-redefs [resume/resume! (fn [ctx] (deliver seen (:llm-config ctx)) {:status :completed})] + (api-control/resume! {:conn c :config config} rid {:model "deepseek-v4-pro"}) + (is (= "deepseek-v4-pro" (:model (deref seen 2000 nil)))))))))) + +(deftest a-resume-leaves-the-supervisor-stream-to-its-stream + ;; The oversight stream opens SUP on the runs row to carry the supervisor's + ;; conversation; a resumed beam rebuilt it as a beam branch and drove it + ;; through the implementor loop. The stream re-opens it itself. + (with-db [c] + (let [rid (runs/start-run! c {:problem "p" :max-turns 10 :beam-width 1})] + (runs/open-branch! c rid {:branch-id "B1"}) + (runs/open-branch! c rid {:branch-id "SUP" :role :supervisor}) + (with-redefs [beam/run-rounds (fn [_ branches _] {:branches branches})] + (is (= ["B1"] (mapv :id (:branches (resume/resume! {:conn c :config {} + :llm-adapter :a :llm-config {} + :run-id rid}))))))))) diff --git a/test/samizdat/infer_test.clj b/test/samizdat/infer_test.clj index 7fb712d3..2cc462f3 100644 --- a/test/samizdat/infer_test.clj +++ b/test/samizdat/infer_test.clj @@ -10,7 +10,9 @@ what a turn would do without running one." (:require [clojure.string :as str] [clojure.test :refer [deftest is testing]] + [samizdat.agent.gates :as gates] [samizdat.agent.infer :as infer] + [samizdat.events :as events] [samizdat.agent.loop :as aloop] [samizdat.agent.state :as state] [samizdat.store.journal :as journal] @@ -452,3 +454,37 @@ (is (= "```tool-call\n" (:prefill b')) "the next request opens inside the fence") (is (= 1 (get-in (state/record-mechanics b {:periodic true}) [:mechanics :periodic])) "and the tally the parse step keeps counts it")))) + + +(deftest a-call-in-a-run-publishes-its-reply-as-it-is-written + ;; The TUI shows the reply while the model writes it: every delta the + ;; client hands on goes onto the bus, tagged with where it starts, so a + ;; reader that missed one can tell rather than garble the rest. + (let [sub (events/subscribe) + threshold gates/threshold] + (try + (with-redefs [gates/threshold (fn [k] (if (= :delta-publish-ms k) 0 (threshold k))) + samizdat.llm.client/chat + (fn [_ _ _ {:keys [on-delta]}] + (on-delta {:reasoning "hm"}) + (on-delta {:text "Hel"}) + (on-delta {:text "lo"}) + {:content "Hello" :finish-reason "stop"})] + ((infer/complete-fn {:llm-adapter ::a :llm-config {} :run-id "R1"}) base-tape)) + (let [ds (filter #(= :delta (:kind %)) (events/collect sub))] + (is (every? #(= ["R1" "B1"] [(:run-id %) (:branch-id %)]) ds)) + (is (= "Hello" (apply str (keep :text ds)))) + (is (= "hm" (apply str (keep :reasoning ds)))) + (is (= [0 3] (keep #(when (:text %) (:text-at %)) ds)) "each piece says where it goes")) + (finally (events/unsubscribe! sub))))) + +(deftest a-call-outside-a-run-publishes-nothing + (let [sub (events/subscribe) + seen (atom :unset)] + (try + (with-redefs [samizdat.llm.client/chat + (fn [_ _ _ opts] (reset! seen (:on-delta opts)) {:content "x" :finish-reason "stop"})] + ((infer/complete-fn {:llm-adapter ::a :llm-config {}}) base-tape)) + (is (nil? @seen) "no watcher asked for, so the call is not streamed") + (is (empty? (filter #(= :delta (:kind %)) (events/collect sub)))) + (finally (events/unsubscribe! sub))))) diff --git a/test/samizdat/llm_stream_test.clj b/test/samizdat/llm_stream_test.clj new file mode 100644 index 00000000..836b7f3e --- /dev/null +++ b/test/samizdat/llm_stream_test.clj @@ -0,0 +1,101 @@ +;; samizdat - a self-hosting agentic harness +;; Copyright (C) 2026 Dmitri Sotnikov +;; +;; This program is free software: you can redistribute it and/or modify +;; it under the terms of the GNU General Public License as published by +;; the Free Software Foundation, either version 3 of the License, or +;; (at your option) any later version. +;; +;; This program is distributed in the hope that it will be useful, +;; but WITHOUT ANY WARRANTY; without even the implied warranty of +;; MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +;; GNU General Public License for more details. +;; +;; You should have received a copy of the GNU General Public License +;; along with this program. If not, see . +;; +;; SPDX-License-Identifier: GPL-3.0-or-later + +(ns samizdat.llm-stream-test + "A provider reply streamed: the chunks folded back into the completion the + adapters already parse, and each delta handed on as it arrives." + (:require [jolt.time] + [clojure.core.async :as async] + [clojure.data.json :as json] + [clojure.test :refer [deftest testing is]] + [ring-chez.adapter :as adapter] + [ring-chez.sse :as rsse] + [samizdat.llm.stream :as stream])) + +(defn- chunk [delta & [extra]] + (merge {:id "c1" :model "m" :choices [{:index 0 :delta delta :finish_reason nil}]} extra)) + +(def ^:private chunks + [(chunk {:role "assistant" :reasoning_content "think"}) + (chunk {:reasoning_content "ing"}) + (chunk {:content "Hel"}) + (chunk {:content "lo"}) + (chunk {:tool_calls [{:index 0 :id "call_1" :type "function" + :function {:name "read_file" :arguments "{\"pa"}}]}) + (chunk {:tool_calls [{:index 0 :function {:arguments "th\":\"a\"}"}}]}) + (merge (chunk {}) {:choices [{:index 0 :delta {} :finish_reason "tool_calls"}]}) + {:id "c1" :model "m" :choices [] :usage {:prompt_tokens 10 :completion_tokens 3 :total_tokens 13}}]) + +(deftest chunks-fold-into-the-completion + (let [c (stream/completion (reduce stream/accumulate nil chunks)) + m (get-in c [:choices 0 :message])] + (is (= "Hello" (:content m))) + (is (= "thinking" (:reasoning_content m))) + (is (= "assistant" (:role m))) + (is (= [{:id "call_1" :type "function" + :function {:name "read_file" :arguments "{\"path\":\"a\"}"}}] + (:tool_calls m))) + (is (= "tool_calls" (get-in c [:choices 0 :finish_reason]))) + (is (= 13 (get-in c [:usage :total_tokens]))) + (is (= "m" (:model c))))) + +(deftest a-delta-is-the-text-and-the-reasoning-it-adds + (is (= {:text "Hel"} (stream/delta (chunk {:content "Hel"})))) + (is (= {:reasoning "think"} (stream/delta (chunk {:reasoning_content "think"})))) + (is (= {:reasoning "r"} (stream/delta (chunk {:reasoning "r"})))) + (is (nil? (stream/delta (chunk {:tool_calls [{:index 0}]}))))) + +(defn- free-port [] (with-open [s (java.net.ServerSocket. 0)] (.getLocalPort s))) + +(defn- with-server [handler f] + (let [port (free-port) + server (adapter/run-server handler {:port port})] + (try (f port) (finally (adapter/stop-server server))))) + +(deftest a-streamed-post-hands-on-each-delta-and-returns-the-whole + (let [seen (atom []) + handler (fn [_] + (let [ch (async/chan 64)] + (future + (doseq [c chunks] + (rsse/send! ch {:data (json/write-str c)}) + (Thread/sleep 5)) + (rsse/send! ch {:data "[DONE]"}) + (async/close! ch)) + {:status 200 :headers {"Content-Type" "text/event-stream"} :body ch}))] + (with-server handler + (fn [port] + (let [r (stream/post (str "http://127.0.0.1:" port "/v1/chat/completions") + {:headers {"Content-Type" "application/json"} :body "{}" + :socket-timeout 5000 :conn-timeout 2000} + #(swap! seen conj %))] + (is (= 200 (:status r))) + (is (= "Hello" (get-in (json/read-str (:body r) :key-fn keyword) + [:choices 0 :message :content]))) + (is (= [{:reasoning "think"} {:reasoning "ing"} {:text "Hel"} {:text "lo"}] @seen) + "each delta, in order, as it came")))))) + +(deftest an-error-status-comes-back-with-its-body + (with-server (fn [_] {:status 400 :headers {"Content-Type" "application/json"} + :body "{\"error\":{\"message\":\"bad request\"}}"}) + (fn [port] + (let [r (stream/post (str "http://127.0.0.1:" port "/x") + {:headers {} :body "{}" :socket-timeout 5000 :conn-timeout 2000} + (fn [_] (throw (ex-info "no deltas on an error" {}))))] + (is (= 400 (:status r))) + (is (= "{\"error\":{\"message\":\"bad request\"}}" (:body r))))))) diff --git a/test/samizdat/llm_test.clj b/test/samizdat/llm_test.clj index a59dd8c9..1638a5ff 100644 --- a/test/samizdat/llm_test.clj +++ b/test/samizdat/llm_test.clj @@ -1809,3 +1809,49 @@ (doseq [h @seen] (is (= "acme" (get h "X-Org"))) (is (= "Bearer k" (get h "Authorization")))))) + +(deftest prose-is-the-reply-without-reasoning-or-call-markup + ;; What a reader of the conversation sees: the words, not the call syntax + ;; the tool chamber already shows, nor reasoning the thinking fold holds. + (testing "a fenced call and a think block" + (is (= "Reading the config first." + (fence/prose (str "which file?Reading the config first.\n" + "```tool-call\n{\"name\": \"read_file\", \"args\": {\"path\": \"a\"}}\n```"))))) + (testing "XML calls, including a pile of stray closers" + (is (= "Let me look." + (fence/prose (str "Let me look.\n\n\n" + "ls\n\n\n" + "\n\n"))))) + (testing "tagged calls, and a reply cut off mid-call or mid-thought" + (is (= "ok" (fence/prose "ok\n{\"name\": \"done\""))) + (is (= "" (fence/prose "still going")))) + (testing "a fence that is not a call is prose" + (is (= "```clojure\n(+ 1 2)\n```" (fence/prose "```clojure\n(+ 1 2)\n```")))) + (testing "blank runs left by the cuts collapse" + (is (= "a\n\nb" (fence/prose "a\n{}\n\n\n\nb"))))) + +(deftest a-call-streams-when-the-endpoint-can-and-someone-is-watching + (let [sent (atom nil) + seen (atom []) + whole "{\"choices\":[{\"message\":{\"content\":\"hi there\"},\"finish_reason\":\"stop\"}]}"] + (with-redefs [samizdat.llm.stream/post (fn [_ {:keys [body]} on-delta] + (reset! sent (json/read-str body :key-fn keyword)) + (on-delta {:text "hi"}) (on-delta {:text " there"}) + {:status 200 :headers {} :body whole}) + jolt.http-client/post (fn [_ _] (throw (ex-info "not streamed" {})))] + (let [r (client/chat (registry/adapter-for :local) + {:base-url "http://h/v1" :model "m" :features #{:stream}} + [{:role "user" :content "x"}] + {:max-tokens 5 :on-delta #(swap! seen conj %)})] + (is (= "hi there" (:content r)) "the reply is the same reply") + (is (true? (:stream @sent))) + (is (= [{:text "hi"} {:text " there"}] @seen))))) + (testing "no watcher, or an endpoint that cannot stream: the whole reply at once" + (doseq [[features opts] [[#{:stream} {:max-tokens 5}] + [#{} {:max-tokens 5 :on-delta identity}]]] + (with-redefs [samizdat.llm.stream/post (fn [& _] (throw (ex-info "streamed" {}))) + jolt.http-client/post (fn [_ _] {:status 200 + :body "{\"choices\":[{\"message\":{\"content\":\"ok\"},\"finish_reason\":\"stop\"}]}"})] + (is (= "ok" (:content (client/chat (registry/adapter-for :local) + {:base-url "http://h/v1" :model "m" :features features} + [{:role "user" :content "x"}] opts)))))))) diff --git a/test/samizdat/test_runner.clj b/test/samizdat/test_runner.clj index 2cc13c0a..bd01521a 100644 --- a/test/samizdat/test_runner.clj +++ b/test/samizdat/test_runner.clj @@ -166,6 +166,7 @@ [samizdat.userspace-adoption-test] [samizdat.event-stream-test] [samizdat.tui-theme-test] + [samizdat.tui-markdown-test] [samizdat.tui-timeline-test] [samizdat.tui-commands-test] [samizdat.live-model-test] @@ -174,6 +175,7 @@ [samizdat.kanban-test] [samizdat.security.policy-test] [samizdat.llm-test] + [samizdat.llm-stream-test] [samizdat.prompt-test] [samizdat.server-test] [samizdat.adapter-test] @@ -272,6 +274,7 @@ samizdat.tui-widgets-test samizdat.store-test samizdat.llm-test + samizdat.llm-stream-test samizdat.agent-test samizdat.approval-test samizdat.base-test @@ -372,6 +375,7 @@ samizdat.userspace-adoption-test samizdat.event-stream-test samizdat.tui-theme-test + samizdat.tui-markdown-test samizdat.tui-timeline-test samizdat.tui-commands-test samizdat.live-model-test diff --git a/test/samizdat/tui_markdown_test.clj b/test/samizdat/tui_markdown_test.clj new file mode 100644 index 00000000..26b8b62b --- /dev/null +++ b/test/samizdat/tui_markdown_test.clj @@ -0,0 +1,39 @@ +;; samizdat - a self-hosting agentic harness +;; Copyright (C) 2026 Dmitri Sotnikov +;; +;; This program is free software: you can redistribute it and/or modify +;; it under the terms of the GNU General Public License as published by +;; the Free Software Foundation, either version 3 of the License, or +;; (at your option) any later version. +;; +;; This program is distributed in the hope that it will be useful, +;; but WITHOUT ANY WARRANTY; without even the implied warranty of +;; MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +;; GNU General Public License for more details. +;; +;; You should have received a copy of the GNU General Public License +;; along with this program. If not, see . +;; +;; SPDX-License-Identifier: GPL-3.0-or-later + +(ns samizdat.tui-markdown-test + "Replies are markdown, drawn the way dirge draws them: headings bold, + bullets as •, code indented and set apart, quotes behind a bar." + (:require [clojure.test :refer [deftest testing is]] + [samizdat.tui.markdown :as md])) + +(defn- shape [lines] (mapv (fn [[_ {:keys [class]} text]] [class text]) lines)) + +(deftest blocks + (is (= [[:md-h1 "Title"] [:md-h2 "Part"] [:md-h2 "Sub"]] + (shape (md/lines "# Title\n## Part\n### Sub")))) + (is (= [[nil " • one"] [nil " • two"] [nil " 3. three"]] + (shape (md/lines "- one\n* two\n3. three")))) + (is (= [[:md-quote "│ said"]] (shape (md/lines "> said")))) + (testing "code keeps its text verbatim, fence lines dropped" + (is (= [[nil "run it:"] [:md-code " (+ 1 2)"] [:md-code " **not bold**"]] + (shape (md/lines "run it:\n```clojure\n(+ 1 2)\n**not bold**\n```"))))) + (testing "inline emphasis markers go, the words stay" + (is (= [[nil "a bold word and `code`"]] (shape (md/lines "a **bold** word and `code`"))))) + (testing "a blank line is a blank row" + (is (= [[nil "a"] [nil ""] [nil "b"]] (shape (md/lines "a\n\nb")))))) diff --git a/test/samizdat/tui_state_test.clj b/test/samizdat/tui_state_test.clj index 81722bde..fe286d7d 100644 --- a/test/samizdat/tui_state_test.clj +++ b/test/samizdat/tui_state_test.clj @@ -269,7 +269,16 @@ (is (= :start (st/enter-action (st/initial "b"))) "nothing to steer, so the words are a problem statement") (is (= :submit (st/enter-action (assoc (st/initial "b") :run-id "r1"))) - "with a run on screen the words are a directive for it")) + "with a run on screen the words are a directive for it") + (is (= :submit (st/enter-action (assoc (st/initial "b") :run-id "r1" + :detail {:run {:status "running"}})))) + (testing "a run that has ended cannot be steered: the words start the next one" + ;; They went out as a directive, the server refused it, and what was typed + ;; was gone — `describe the project` typed over a finished run. + (doseq [status ["completed" "failed" "aborted" "exhausted"]] + (is (= :start (st/enter-action (assoc (st/initial "b") :run-id "r1" + :detail {:run {:status status}}))) + status)))) (deftest a-started-run-is-selected-and-the-box-is-emptied (let [s (-> (st/initial "b") @@ -357,3 +366,32 @@ (deftest a-newline-can-be-typed-on-purpose (is (= "one\n" (:input (st/newline (st/set-input (st/initial "b") "one")))))) + +(deftest a-reply-being-written-is-folded-in-as-it-arrives + (let [ev (fn [data] {:event "delta" :data (merge {:branch_id "B1"} data)}) + fold (fn [s e] (first (st/apply-event s e))) + s (reduce fold (st/initial "b") + [(ev {:reasoning "hm" :reasoning-at 0}) + (ev {:text "Hel" :text-at 0}) + (ev {:text "lo" :text-at 3})])] + (is (= {:text "Hello" :reasoning "hm"} (get-in s [:live "B1"]))) + (is (= #{} (second (st/apply-event s (ev {:text "!" :text-at 5})))) + "nothing to fetch: the event carries it") + (testing "a piece already held is not doubled" + (is (= "Hello" (get-in (fold s (ev {:text "lo" :text-at 3})) [:live "B1" :text])))) + (testing "the branch's next turn row takes over from it" + (is (nil? (get-in (fold s {:id "9" :event "turn" :data {:branch_id "B1"}}) [:live "B1"])))))) + +(deftest a-pushed-event-is-read-as-it-comes-off-the-wire + ;; The stream hands over each event's data as the JSON TEXT it was sent + ;; as. Read as a map it never was, every event lost its branch: a turn on + ;; the branch on screen asked for no branch fetch — the conversation only + ;; caught up when something else refreshed it — and a reply being written + ;; was filed under no branch at all. + (let [s (assoc (st/initial "b") :run-id "R" :branch-id "B1")] + (is (contains? (second (st/apply-event s {:id "41" :event "turn" + :data "{\"id\":41,\"branch_id\":\"B1\",\"kind\":\"turn\"}"})) + :branch)) + (is (= "Hel" (get-in (first (st/apply-event s {:event "delta" + :data "{\"branch_id\":\"B1\",\"text\":\"Hel\",\"text-at\":0}"})) + [:live "B1" :text]))))) diff --git a/test/samizdat/tui_timeline_test.clj b/test/samizdat/tui_timeline_test.clj index 1e3e2417..352bc7c6 100644 --- a/test/samizdat/tui_timeline_test.clj +++ b/test/samizdat/tui_timeline_test.clj @@ -53,7 +53,7 @@ (deftest everyone-speaks-in-the-order-it-happened (let [es (tl/entries state settings)] - (is (= [[:user :say] [:agent :say] [:agent :thinking] [:agent :tool] + (is (= [[:user :say] [:agent :thinking] [:agent :say] [:agent :tool] [:user :say] [:critic :say] [:agent :say] [:agent :tool] [:supervisor :say] [:supervisor :say]] (mapv (juxt :role :kind) es))) @@ -86,21 +86,24 @@ (is (= before (take (count before) after))))) ;; --- following the bottom ------------------------------------------------------- +;; +;; A pane scrolls by ROWS (ftxui's :scroll): :top is the first row shown, nil +;; follows the bottom. The pane reports {:top :max :rows}; keys page from it. -(deftest the-view-follows-the-bottom-until-you-scroll-up - (let [ks (mapv str (range 30)) - s (st/initial "b")] - (is (nil? (:scroll-anchor s)) "following") - (let [up (st/scroll s ks -10)] - (is (= "19" (:scroll-anchor up)) "ten entries up from the bottom") - (is (= "9" (:scroll-anchor (st/scroll up ks -10)))) - (is (= "0" (:scroll-anchor (st/scroll up ks -100))) "not past the top") - (testing "scrolling back down to the bottom follows again" - (is (nil? (:scroll-anchor (st/scroll up ks 10)))) - (is (nil? (:scroll-anchor (st/scroll up ks 100))))) - (testing "an anchor holds its place as entries arrive below it" - (is (= "19" (:scroll-anchor (st/scroll up (conj ks "30" "31") 0))))) - (is (nil? (:scroll-anchor (st/follow up))))))) +(deftest a-pane-follows-the-bottom-until-paged-up + (let [s (st/scrolled (st/initial "b") :conversation {:top nil :max 100 :rows 22})] + (is (nil? (st/scroll-top s :conversation)) "following") + (let [up (st/page s :conversation -1)] + (is (= 80 (st/scroll-top up :conversation)) "a page is the rows shown, less two") + (is (= 60 (st/scroll-top (st/page up :conversation -1) :conversation))) + (is (= 0 (st/scroll-top (reduce #(st/page %1 :conversation %2) up (repeat 9 -1)) + :conversation)) + "not past the top") + (testing "paging back down to the end follows again" + (is (nil? (st/scroll-top (st/page up :conversation 1) :conversation)))) + (is (nil? (st/scroll-top (st/follow up :conversation) :conversation)))) + (testing "what the pane reports is what is kept" + (is (= 12 (st/scroll-top (st/scrolled s :activity {:top 12 :max 40 :rows 10}) :activity)))))) (deftest only-folds-that-are-drawn-are-folds (let [es (tl/entries (assoc-in state [:branch :turns 1 :result] "1\n2\n3\n4\n5") settings)] @@ -114,3 +117,38 @@ (is (= #{"t2/tool/more"} (:expanded opened)) "the newest collapsed thing opens") (is (= #{} (:expanded (st/toggle-latest-fold opened ids))) "pressed again, it shuts") (is (= s (st/toggle-latest-fold s [])) "nothing to open, nothing changes"))) + +(deftest the-agents-words-are-drawn-without-markup + ;; A reply's call syntax is the chamber's to show and its the + ;; thinking fold's; drawn raw they were a column of `` lines. + (let [s {:run-id "R" :branch-id "B1" + :branch {:turns [{:turn 1 :tool_name "shell" :args "{\"command\":\"ls\"}" + :result "a" :category "success" :created_at "t1"}]} + :turn-text {1 {:assistant_text "plan itListing.\n\n" + :reasoning_text "first"}}} + es (tl/entries s {}) + say (first (filter #(= :say (:kind %)) es)) + thinking (first (filter #(= :thinking (:kind %)) es))] + (is (= "Listing." (:text say))) + (is (= "first\nplan it" (:text thinking)) "inline reasoning joins the thinking fold")) + (testing "a turn that said nothing but a call has no say entry" + (let [s {:branch {:turns [{:turn 1 :tool_name "shell" :created_at "t1"}]} + :turn-text {1 {:assistant_text "```tool-call\n{\"name\": \"shell\"}\n```"}}}] + (is (not-any? #(= :say (:kind %)) (tl/entries s {})))))) + +(deftest a-notes-reasoning-folds-too + (let [s {:branch {:notes [{:id 7 :kind "critic-score" :created_at "t1" + :data "{\"reply\":\"hmmProgress 2.\"}"}]}} + es (tl/entries s settings)] + (is (= ["Progress 2."] (map :text (filter #(= :say (:kind %)) es)))) + (is (= ["hmm"] (map :text (filter #(= :thinking (:kind %)) es)))))) + +(deftest the-reply-being-written-is-the-last-entry + (let [s {:branch-id "B1" + :branch {:turns [{:turn 1 :tool_name "shell" :created_at "t1"}]} + :live {"B1" {:text "soWriting it\n" :reasoning "first"} + "B2" {:text "not this branch"}}} + es (tl/entries s {})] + (is (= [[:agent :thinking "first\nso"] [:agent :say "Writing it"]] + (mapv (juxt :role :kind :text) (take-last 2 es)))) + (is (every? :live? (take-last 2 es))))) diff --git a/test/samizdat/tui_widgets_test.clj b/test/samizdat/tui_widgets_test.clj index f1c91701..f5e43742 100644 --- a/test/samizdat/tui_widgets_test.clj +++ b/test/samizdat/tui_widgets_test.clj @@ -191,14 +191,13 @@ out (conv {:branch {:turns turns}} {:turns 40})] (is (= 3 (count (keep #(:key (props-of %)) (nodes-of :vbox out)))))))) -(deftest the-newest-entry-holds-the-focus-unless-scrolled-away - ;; ftxui scrolls a frame to its focused element: that is how the pane - ;; follows the bottom as the run speaks. - (let [focused #(keep (fn [n] (when (:focus (props-of n)) (:key (props-of n)))) - (nodes-of :vbox %))] - (is (= ["t2/tool"] (focused (conv {:branch {:turns turns}} {})))) - (is (= ["t1/tool"] (focused (conv {:branch {:turns turns} :scroll-anchor "t1/tool"} {}))) - "scrolled up, the anchor holds the view"))) +(deftest the-conversation-is-a-scroll-pane-at-the-top-it-was-left + (let [pane #(first (nodes-of :scroll %))] + (is (some? (pane (conv {:branch {:turns turns}} {})))) + (is (nil? (:top (props-of (pane (conv {:branch {:turns turns}} {}))))) "following") + (is (= 7 (:top (props-of (pane (conv {:branch {:turns turns} + :scroll {:conversation {:top 7}}} {}))))) + "scrolled up, it stays there"))) (deftest an-expanded-fold-is-open-and-its-body-is-drawn (let [id (w/fold-id 1 :result) diff --git a/tui/samizdat/tui/core.clj b/tui/samizdat/tui/core.clj index e7146e25..dc638a2e 100644 --- a/tui/samizdat/tui/core.clj +++ b/tui/samizdat/tui/core.clj @@ -206,7 +206,7 @@ :help (apply say! "commands:" (cmd/help-lines cs)) :quit (ui/exit!) :clear (swap! state st/clear-local) - :follow (swap! state st/follow) + :follow (swap! state st/follow :conversation) :abort (abort!) :resume (resume!) :start (if arg (start! arg) (say! "/run needs the problem to work on")) @@ -276,7 +276,8 @@ (cmd/parse text {}) (command! text) - :else (do (swap! state st/remember-input text) + ;; Sending is looking at the bottom again, as in dirge. + :else (do (swap! state #(-> % (st/remember-input text) (st/follow :conversation))) (action text))))) (def ^:private handlers @@ -293,7 +294,8 @@ :abort abort! :resume resume! :reply (fn [kind id] (swap! state st/start-reply kind id)) - :toggle-option #(swap! state st/toggle-option %)}) + :toggle-option #(swap! state st/toggle-option %) + :scroll (fn [pane view] (swap! state st/scrolled pane view))}) ;; --- the feeds ---------------------------------------------------------------- ;; @@ -460,9 +462,6 @@ ;; along with the ones the layout names — one stylesheet for both. (theme/apply-theme th (layout/expand layout s)))) -(def ^:private page-entries 10) -(def ^:private wheel-entries 3) - (defn- conversation-entries "The conversation's entries as the widget draws them, for the keys that move through it." @@ -478,10 +477,8 @@ key has to mean it. `state/pending-decision` decides whether there is such a dialog — over a questionnaire's answer box a `y` is a letter being typed, and it goes through untouched." - [{:keys [type key char control button]}] - (let [scroll! (fn [delta] - (swap! state #(st/scroll % (mapv :key (conversation-entries %)) delta)) - true)] + [{:keys [type key char control]}] + (let [page! (fn [dir] (swap! state st/page :conversation dir) true)] (cond (and (= :key type) (= :ctrl-c key)) (do (ui/exit!) true) ;; ftxui delivers Ctrl+Q as a :key, never as a :character with a @@ -491,14 +488,13 @@ (and (= :key type) (= :f5 key)) (do (future (poll-once!)) true) ;; The conversation: page through it, and back to following the bottom. - (and (= :key type) (= :page-up key)) (scroll! (- page-entries)) - (and (= :key type) (= :page-down key)) (scroll! page-entries) - (and (= :mouse type) (= :wheel-up button)) (scroll! (- wheel-entries)) - (and (= :mouse type) (= :wheel-down button)) (scroll! wheel-entries) + ;; The wheel is not handled here: each pane scrolls under the mouse. + (and (= :key type) (= :page-up key)) (page! -1) + (and (= :key type) (= :page-down key)) (page! 1) ;; Down or End while scrolled up jumps back to the bottom (dirge); ;; while following, they are the input's. - (and (= :key type) (#{:arrow-down :end} key) (:scroll-anchor @state)) - (do (swap! state st/follow) true) + (and (= :key type) (#{:arrow-down :end} key) (st/scroll-top @state :conversation)) + (do (swap! state st/follow :conversation) true) ;; Tab completes a slash command's name. (and (= :key type) (= :tab key) (str/starts-with? (str (:input @state)) "/")) (do (swap! state #(st/set-input % (cmd/complete (:input %) (:commands (layout/current))))) diff --git a/tui/samizdat/tui/markdown.clj b/tui/samizdat/tui/markdown.clj new file mode 100644 index 00000000..73d6e44d --- /dev/null +++ b/tui/samizdat/tui/markdown.clj @@ -0,0 +1,52 @@ +;; samizdat - a self-hosting agentic harness +;; Copyright (C) 2026 Dmitri Sotnikov +;; +;; This program is free software: you can redistribute it and/or modify +;; it under the terms of the GNU General Public License as published by +;; the Free Software Foundation, either version 3 of the License, or +;; (at your option) any later version. +;; +;; This program is distributed in the hope that it will be useful, +;; but WITHOUT ANY WARRANTY; without even the implied warranty of +;; MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +;; GNU General Public License for more details. +;; +;; You should have received a copy of the GNU General Public License +;; along with this program. If not, see . +;; +;; SPDX-License-Identifier: GPL-3.0-or-later + +(ns samizdat.tui.markdown + "A reply's markdown as rows, the way dirge draws it: headings bold, bullets + as •, code indented and set apart, quotes behind a bar. Block structure + only — inline emphasis markers are dropped, since a styled span inside a + wrapped line would stop it wrapping. Pure: hiccup rows for the widgets." + (:require [clojure.string :as str])) + +(defn- row + ([text] (row nil text)) + ([class text] [:wrapped (cond-> {} class (assoc :class class)) text])) + +(defn- inline [s] + (-> s + (str/replace #"\*\*(.+?)\*\*" "$1") + (str/replace #"__(.+?)__" "$1"))) + +(defn- block-row [line] + (condp re-find line + #"^#\s+(.*)" :>> #(row :md-h1 (inline (second %))) + #"^#{2,6}\s+(.*)" :>> #(row :md-h2 (inline (second %))) + #"^\s*[-*+]\s+(.*)" :>> #(row (str " • " (inline (second %)))) + #"^\s*(\d+)[.)]\s+(.*)" :>> #(row (str " " (nth % 1) ". " (inline (nth % 2)))) + #"^>\s?(.*)" :>> #(row :md-quote (str "│ " (inline (second %)))) + (row (inline line)))) + +(defn lines + "`text` as a vector of rows." + [text] + (loop [[l & more :as ls] (str/split-lines (str text)) code? false out []] + (cond + (empty? ls) out + (str/starts-with? (str/triml l) "```") (recur more (not code?) out) + code? (recur more code? (conj out (row :md-code (str " " l)))) + :else (recur more code? (conj out (block-row l)))))) diff --git a/tui/samizdat/tui/state.clj b/tui/samizdat/tui/state.clj index 203cc615..91b5a0af 100644 --- a/tui/samizdat/tui/state.clj +++ b/tui/samizdat/tui/state.clj @@ -27,7 +27,8 @@ and the layout never has to name data, only widgets; that is what lets a user put a panel anywhere without anything being rewired." (:refer-clojure :exclude [newline]) - (:require [clojure.string :as str])) + (:require [clojure.string :as str] + [samizdat.api.sse :as sse])) (def handler-keys "Every action a widget may ask the loop to take. @@ -39,7 +40,7 @@ that called them, which no test could see because each half was correct on its own." #{:decide :answer :toggle :select-run :select-branch :input :submit :start - :abort :resume :reply :toggle-option}) + :abort :resume :reply :toggle-option :scroll}) (def max-trace "How many steps the UI holds. The server's ring is bounded and so is this: @@ -73,9 +74,8 @@ ;; Whether the run's event stream is up. While it is, the run panels are ;; refreshed when an event says they changed rather than on a timer. :live? false - ;; Where the conversation is scrolled to: the key of the entry held in - ;; view, or nil to follow the bottom as new entries arrive. - :scroll-anchor nil + ;; Each pane's scroll, as it last reported it: {pane {:top :max :rows}}. + :scroll {} ;; What this TUI printed — command output, /help — drawn in the ;; conversation as the harness's voice. Local: nothing the server holds. :local-notes [] @@ -363,9 +363,14 @@ A pure function rather than a branch inside the widget so the rule is testable, and named as a KEY so the widget still spells both handlers out literally — `every-handler-the-loop-offers-has-a-caller` reads those off - the source." + the source. + + A run that has ENDED has nothing to steer either: a directive to it is + refused, and what was typed was lost. Its status is read off the run + detail; until that has loaded, a selected run is taken as live." [s] - (if (:run-id s) :submit :start)) + (let [status (some-> (get-in s [:detail :run :status]) str)] + (if (and (:run-id s) (or (nil? status) (= "running" status))) :submit :start))) (defn apply-start "Fold the answer to POST /v1/runs. @@ -405,12 +410,33 @@ [d] (select-keys d [:node :cell :transition :ms :failed :turn :branch_id])) +(defn- splice + "`piece` into `held` at offset `at`: appended when it follows on, laid + over what is held when it repeats some of it, and after the gap when one + was lost — this is a preview the turn row replaces, not the record." + [held piece at] + (let [held (str held) + at (or at (count held))] + (str (subs held 0 (min at (count held))) piece))) + +(defn- live-delta + "A piece of the reply a branch is writing (samizdat.agent.infer publishes + them while a model call streams)." + [s {:keys [branch_id text reasoning] :as d}] + (update-in s [:live branch_id] + (fn [l] + (cond-> (or l {:text "" :reasoning ""}) + text (update :text splice text (:text-at d)) + reasoning (update :reasoning splice reasoning (:reasoning-at d)))))) + (defn apply-event "Fold one pushed event into the state. Returns [state wants]: `wants` is the set of things the event changed that the stream does not carry — :detail, :branch, :approvals, :runs — for the caller to fetch, coalesced." [s {:keys [id event data]}] - (let [s (cond-> s id (assoc :journal-cursor (or (parse-long (str id)) (:journal-cursor s))))] + (let [;; Off the wire an event's data is the JSON text it was sent as. + data (if (string? data) (:data (sse/parse-data data)) data) + s (cond-> s id (assoc :journal-cursor (or (parse-long (str id)) (:journal-cursor s))))] (case event "step" [(update s :trace (fn [t] (let [t (conj (vec t) (step-entry data))] (if (> (count t) max-trace) @@ -418,8 +444,11 @@ t)))) #{}] "approval" [s #{:approvals}] - [s (cond-> (get refresh-for event #{:detail}) - (and (:branch_id data) (= (:branch_id data) (:branch-id s))) (conj :branch))]))) + "delta" [(live-delta s data) #{}] + ;; A branch's turn row is in: it holds what the live preview showed. + [(cond-> s (= "turn" event) (update :live dissoc (:branch_id data))) + (cond-> (get refresh-for event #{:detail}) + (and (:branch_id data) (= (:branch_id data) (:branch-id s))) (conj :branch))]))) (defn stream-status "Note the event stream's state: `status` 200 is up, anything else down." @@ -427,26 +456,30 @@ (assoc s :live? (= 200 status))) ;; --- following the bottom ---------------------------------------------------- +;; +;; A pane scrolls by rows (ftxui's :scroll). What it last reported is kept per +;; pane, {:top :max :rows}: :top the first row shown, nil following the bottom. + +(defn scrolled + "Keep what pane `pane` reported: its view and how far it can go." + [s pane view] + (assoc-in s [:scroll pane] view)) + +(defn scroll-top [s pane] (get-in s [:scroll pane :top])) (defn follow - "Follow the bottom of the conversation again." - [s] - (assoc s :scroll-anchor nil)) - -(defn scroll - "Move the conversation `delta` entries (negative is up) over `ks`, the - entries' keys in order. Following the bottom is the anchor being nil; - scrolling back down to the last entry follows again, and an anchor holds - its entry however many arrive below it — the view stays where the reader - left it (dirge's rule)." - [s ks delta] - (let [n (count ks) - at (or (some-> (:scroll-anchor s) (#(.indexOf ^java.util.List ks %)) (#(when (>= % 0) %))) - (dec n)) - to (max 0 (min (dec n) (+ at delta)))] - (if (or (zero? n) (>= to (dec n))) - (follow s) - (assoc s :scroll-anchor (nth ks to))))) + "Follow the bottom of `pane` again." + [s pane] + (assoc-in s [:scroll pane :top] nil)) + +(defn page + "Move `pane` a page up (`dir` -1) or down (1): the rows it shows, less two + for context. Paging back to the end follows again, as dirge does." + [s pane dir] + (let [{:keys [top max rows]} (get-in s [:scroll pane]) + max (or max 0) + to (clojure.core/max 0 (+ (or top max) (* dir (clojure.core/max 1 (- (or rows 3) 2)))))] + (assoc-in s [:scroll pane :top] (when (< to max) to)))) (defn toggle-latest-fold "Ctrl+O, from dirge: open the newest fold in `ids` (in order), or shut it diff --git a/tui/samizdat/tui/timeline.clj b/tui/samizdat/tui/timeline.clj index 85c40550..969247e7 100644 --- a/tui/samizdat/tui/timeline.clj +++ b/tui/samizdat/tui/timeline.clj @@ -32,7 +32,8 @@ (:require ;; the java.time.* host shim, before data.json — see samizdat.store.journal [jolt.time] [clojure.data.json :as json] - [clojure.string :as str])) + [clojure.string :as str] + [samizdat.llm.fence :as fence])) (defn- parse [s] (if (string? s) @@ -52,18 +53,26 @@ (defn- failed? [t] (contains? #{"failure" "mechanics"} (str (:category t)))) +(defn- words + "`s` split into what is said and what was thought: the prose with call + markup and reasoning blocks removed (fence/prose), and the reasoning those + blocks held joined onto `reasoning`." + [s reasoning] + {:say (fence/prose s) + :thinking (str/join "\n" (remove str/blank? [(str reasoning) (fence/reasoning-of s)]))}) + (defn- turn-entries [t text banner] (let [n (:turn t) at (str (:created_at t)) - {:keys [assistant_text reasoning_text]} text] + {:keys [say thinking]} (words (:assistant_text text) (:reasoning_text text))] (cond-> [] - (not (str/blank? (str assistant_text))) - (conj {:key (str "t" n "/say") :role :agent :kind :say :turn n :at at - :text (str assistant_text)}) - - (not (str/blank? (str reasoning_text))) + (not (str/blank? thinking)) (conj {:key (str "t" n "/thinking") :role :agent :kind :thinking :turn n :at at - :text (str reasoning_text)}) + :text thinking}) + + (not (str/blank? say)) + (conj {:key (str "t" n "/say") :role :agent :kind :say :turn n :at at + :text say}) :always (conj {:key (str "t" n "/tool") :role :agent :kind :tool :turn n :at at @@ -95,9 +104,26 @@ (defn- note-entries [notes by-kind] (for [n notes :let [{:keys [role say]} (get by-kind (str (:kind n)))] - :when role] - {:key (str "n" (:id n)) :role role :kind :say :at (str (:created_at n)) - :text (note-text (parse (:data n)) say) :note (str (:kind n))})) + :when role + :let [w (words (note-text (parse (:data n)) say) nil) + base {:role role :at (str (:created_at n)) :note (str (:kind n))}] + e [(when-not (str/blank? (:thinking w)) + (assoc base :key (str "n" (:id n) "/thinking") :kind :thinking :text (:thinking w))) + (when-not (str/blank? (:say w)) + (assoc base :key (str "n" (:id n)) :kind :say :text (:say w)))] + :when e] + e)) + +(defn- live-entries + "The reply the branch is writing right now, as it streams in: after + everything else, marked :live?, and gone when its turn row lands." + [{:keys [text reasoning]}] + (let [{:keys [say thinking]} (words text reasoning)] + (cond-> [] + (not (str/blank? thinking)) + (conj {:key "live/thinking" :role :agent :kind :thinking :text thinking :live? true}) + (not (str/blank? say)) + (conj {:key "live/say" :role :agent :kind :say :text say :live? true})))) (defn entries "The branch on screen as a vector of entries, oldest first: @@ -118,9 +144,11 @@ (for [n (:local-notes state)] {:key (:key n) :role :system :kind :say :at (:at n) :text (:text n)})))] - (vec (cond->> later - (not (str/blank? (str problem))) - (cons {:key "problem" :role :user :kind :say :at "" :text (str problem)}))))) + (vec (concat + (cond->> later + (not (str/blank? (str problem))) + (cons {:key "problem" :role :user :kind :say :at "" :text (str problem)})) + (live-entries (get-in state [:live (:branch-id state)])))))) (defn fold-id "The id of the fold an entry's body would sit behind." diff --git a/tui/samizdat/tui/widgets.clj b/tui/samizdat/tui/widgets.clj index 6b97d598..27ce163e 100644 --- a/tui/samizdat/tui/widgets.clj +++ b/tui/samizdat/tui/widgets.clj @@ -39,6 +39,7 @@ [samizdat.tui.layout :as layout] [samizdat.tui.state :as st] [samizdat.tui.commands :as cmd] + [samizdat.tui.markdown :as md] [samizdat.tui.timeline :as tl])) ;; --- shared shapes ----------------------------------------------------------- @@ -139,6 +140,18 @@ :on-change (toggle-fn state id)} (if open? (body-fn) [:empty])])) +;; --- scrolling ---------------------------------------------------------------- + +(defn- pane + "`content` in a pane that scrolls by rows: the wheel over it scrolls it, + and it follows the bottom until it is scrolled up. Where it is lives in the + state under `id`, so keys can page it (state/page)." + [state id content] + [:scroll {:flex true + :top (get-in state [:scroll id :top]) + :on-change (fn [view] (when-let [f (get-in state [:on :scroll])] (f id view)))} + content]) + ;; --- the activity log -------------------------------------------------------- (defn- step-line [{:keys [node cell transition ms failed]}] @@ -168,7 +181,7 @@ (when (and dropped (pos? dropped)) [:text {:class :warn} (str " ⚠ " dropped " step(s) dropped")]) (if (seq trace) - (into [:vbox {:flex true}] (map step-line trace)) + (pane state :activity (into [:vbox] (map step-line trace))) (empty-note "no steps yet"))))) ;; --- the conversation -------------------------------------------------------- @@ -208,14 +221,15 @@ (def ^:private writing-tools #{"write_file" "edit_file" "patch"}) (defn- say-entry - "A line in a role's voice, dirge-style: ` ` and the words wrapped - under it with a hanging indent. `:wrapped` rather than `:paragraph`: a - path or a token wider than the pane breaks instead of being cut off." + "A line in a role's voice, dirge-style: ` ` and the words under it + with a hanging indent, drawn as markdown (samizdat.tui.markdown) in the + role's colour. Rows are `:wrapped`: a path or a token wider than the pane + breaks instead of being cut off." [{:keys [role text pending?]} handles] [:hbox [:text {:class role} (str "<" (get handles role (name role)) "> ")] - [:wrapped {:class role :flex true} - (str text (when pending? " (not applied yet)"))]]) + (into [:vbox {:class role :flex true}] + (md/lines (str text (when pending? " (not applied yet)"))))]) (defn- thinking-entry [state e] (fold state (tl/fold-id e) @@ -264,27 +278,26 @@ A tool call is a chamber showing the first lines of its result; the rest, and the model's thinking, fold, closed by default. - FOLLOWS THE BOTTOM. The newest entry holds the frame's focus, so the pane - scrolls as the run speaks; scrolling up anchors it on an entry and it - stays there as more arrive, until End (or scrolling back down) follows - again. + FOLLOWS THE BOTTOM. A scroll pane (`pane`): the wheel over it or PgUp + scrolls it by rows and it stays there as more arrive, until End (or + scrolling back down) follows again. BOUNDED: `:turns` in the layout, the newest that many." [state props] (let [cfg (get-in state [:settings :conversation]) es (shown-entries (tl/entries state cfg) (or (:turns props) default-turns-shown)) - focus (or (:scroll-anchor state) (:key (peek es))) handles (:handles cfg)] (panel (assoc props :title (or (:title props) (some->> (:branch-id state) (str "BRANCH ")))) (if (seq es) - (into [:vbox {:flex true :frame :y :scroll-indicator :v}] - (for [e es] - [:vbox (cond-> {:key (:key e)} (= focus (:key e)) (assoc :focus true)) - (case (:kind e) - :say (say-entry e handles) - :thinking (thinking-entry state e) - :tool (tool-entry state e cfg))])) + (pane state :conversation + (into [:vbox] + (for [e es] + [:vbox {:key (:key e)} + (case (:kind e) + :say (say-entry e handles) + :thinking (thinking-entry state e) + :tool (tool-entry state e cfg))]))) (empty-note "no turns yet — pick a run"))))) ;; --- the side panels --------------------------------------------------------- @@ -699,6 +712,9 @@ rows [:<= (or (:max-lines props) 8)] row [:hbox (merge (select-keys props [:flex :width]) {:height rows}) [:input {:flex true + ;; Typing into a TUI that has just opened lands here, + ;; not in whichever widget the layout reaches first. + :autofocus true :multiline true :wrap true :value (or (:input state) "")