Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 16 additions & 8 deletions docs/tui.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down Expand Up @@ -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:

- `<you>` — the problem, and every steer a person sent;
- `<agent>` — what the model said; its thinking folds under `◇ thinking`;
- `<agent>` — 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 `<think>` 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;
Expand All @@ -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

Expand Down Expand Up @@ -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
Expand Down
10 changes: 10 additions & 0 deletions resources/gates.edn
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
14 changes: 11 additions & 3 deletions resources/tui.edn
Original file line number Diff line number Diff line change
Expand Up @@ -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"}
Expand Down Expand Up @@ -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 {}]]}
5 changes: 4 additions & 1 deletion src/samizdat/agent/beam.clj
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
42 changes: 39 additions & 3 deletions src/samizdat/agent/infer.clj
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down Expand Up @@ -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}).

Expand Down Expand Up @@ -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).
Expand All @@ -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`
Expand All @@ -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
Expand Down
4 changes: 3 additions & 1 deletion src/samizdat/agent/resume.clj
Original file line number Diff line number Diff line change
Expand Up @@ -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)))]
Expand Down
34 changes: 28 additions & 6 deletions src/samizdat/api/control.clj
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
59 changes: 47 additions & 12 deletions src/samizdat/api/sse.clj
Original file line number Diff line number Diff line change
Expand Up @@ -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"))

Expand Down Expand Up @@ -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]
Expand All @@ -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
Expand All @@ -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 []}))
Expand All @@ -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"
Expand Down
Loading
Loading