From 8ca7b90cca53cbf26bfb70a284deb8a07ff6cea5 Mon Sep 17 00:00:00 2001 From: Yogthos Date: Wed, 23 Sep 2026 23:44:32 -0400 Subject: [PATCH 1/2] Refuse to resume a run that is still running, and resume on its own model A resume of a live run started a second driver over the same branches, writing duplicate turns, and it took the config's default model rather than the run's, so every one of its calls went to an endpoint the run never used. The active slot is now claimed atomically before the driver spawns, the run row records the provider alias it was started on, and a resume with no model in the request picks that back up. The supervisor stream's branch is left for the stream to re-open instead of being driven as a beam branch. --- src/samizdat/agent/beam.clj | 5 ++- src/samizdat/agent/resume.clj | 4 +- src/samizdat/api/control.clj | 34 ++++++++++++++--- src/samizdat/workflow.clj | 5 ++- test/samizdat/control_test.clj | 68 +++++++++++++++++++++++++++++----- 5 files changed, 98 insertions(+), 18 deletions(-) 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/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..cf7f4bca 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 already running in this process") + (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/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/samizdat/control_test.clj b/test/samizdat/control_test.clj index 822b41e7..f5d750ce 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])) "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}))))))))) From d0b9d41056ca90359314128a6f1819615421e31f Mon Sep 17 00:00:00 2001 From: Yogthos Date: Thu, 24 Sep 2026 05:49:47 -0400 Subject: [PATCH 2/2] Stream replies into the TUI, render the agent log cleanly, scroll panes by row Model calls stream when the endpoint supports it (:stream feature) and a run is watching: samizdat.llm.stream speaks HTTP/1.1 over the socket, hands each delta on, and folds the chunks back into the completion the adapters already parse. infer publishes the deltas on the event bus, throttled by gates.edn :delta-publish-ms, and the TUI draws the reply as it is written until the turn row replaces it. The TUI never parsed the JSON of pushed events, so a turn on the branch on screen never asked for a branch refresh and the conversation lagged until something else refreshed it. Fixed, which also feeds the activity pane from pushes. The agent log drops call markup and folds blocks (fence/prose), renders replies as light markdown, and scrolls by rows through ftxui-jolt's :scroll pane: the wheel scrolls the pane under the pointer, PgUp/PgDn page the conversation, and the activity log scrolls too. The compose box has the focus at startup, runs to the right edge, and Enter starts a new run when the selected one has ended instead of sending a directive it would refuse. --- docs/tui.md | 24 ++-- resources/gates.edn | 10 ++ resources/tui.edn | 14 ++- src/samizdat/agent/infer.clj | 42 ++++++- src/samizdat/api/control.clj | 2 +- src/samizdat/api/sse.clj | 59 +++++++-- src/samizdat/config.clj | 10 +- src/samizdat/llm/client.clj | 31 +++-- src/samizdat/llm/fence.clj | 23 ++++ src/samizdat/llm/stream.clj | 173 +++++++++++++++++++++++++++ test-tui/samizdat/tui/mouse_test.clj | 62 +++++++++- test/samizdat/base_test.clj | 3 +- test/samizdat/control_test.clj | 2 +- test/samizdat/infer_test.clj | 36 ++++++ test/samizdat/llm_stream_test.clj | 101 ++++++++++++++++ test/samizdat/llm_test.clj | 46 +++++++ test/samizdat/test_runner.clj | 4 + test/samizdat/tui_markdown_test.clj | 39 ++++++ test/samizdat/tui_state_test.clj | 40 ++++++- test/samizdat/tui_timeline_test.clj | 68 ++++++++--- test/samizdat/tui_widgets_test.clj | 15 ++- tui/samizdat/tui/core.clj | 28 ++--- tui/samizdat/tui/markdown.clj | 52 ++++++++ tui/samizdat/tui/state.clj | 89 +++++++++----- tui/samizdat/tui/timeline.clj | 56 ++++++--- tui/samizdat/tui/widgets.clj | 52 +++++--- 26 files changed, 931 insertions(+), 150 deletions(-) create mode 100644 src/samizdat/llm/stream.clj create mode 100644 test/samizdat/llm_stream_test.clj create mode 100644 test/samizdat/tui_markdown_test.clj create mode 100644 tui/samizdat/tui/markdown.clj 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/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/api/control.clj b/src/samizdat/api/control.clj index cf7f4bca..dcb2b341 100644 --- a/src/samizdat/api/control.clj +++ b/src/samizdat/api/control.clj @@ -252,7 +252,7 @@ ;; 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 already running in this process") + (refuse "is still running") (do (let [cancel* (atom nil) started (cancel/start! 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/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 f5d750ce..7c5b4795 100644 --- a/test/samizdat/control_test.clj +++ b/test/samizdat/control_test.clj @@ -891,7 +891,7 @@ (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])) "running"))) + (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)))))) 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) "")