-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathagent.py
More file actions
1392 lines (1197 loc) · 61.7 KB
/
Copy pathagent.py
File metadata and controls
1392 lines (1197 loc) · 61.7 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
"""Grad -- the agent loop (HANDOFF §3, §9, §12 step 1).
A `ClaudeSDKClient` multi-turn session with a small system prompt, the six
built-in tools, a deny-by-default permission mode, and a `PreToolUse` gate. The
custom capability is not here: it is the CLIs in `tools/`, reached over Bash.
Three configuration details are load-bearing and easy to get wrong, so they are
asserted rather than assumed:
* `allowed_tools` is an *auto-approve* list, not a sandbox. Built-in tools stay
in the model's toolset regardless of what is listed, so the restriction comes
from `disallowed_tools` (deny rules beat every other step) plus the mode.
* the permission mode's name and semantics have changed between SDK releases,
so `agent.py probe` attempts a call that should be denied and reports whether
it was *denied*, not prompted and not silently allowed. Re-run it after any
SDK upgrade.
* `setting_sources` is left unset, so a stray `settings.json` cannot add allow
rules silently. The whole permission configuration lives in code.
"""
from __future__ import annotations
import argparse
import asyncio
import dataclasses
import json
import logging
import os
import sys
import time
from pathlib import Path
from typing import Any
import hooks
from core import appdata, config as config_mod, credentials, effort, migrate, paths, quota_log
from core.errors import EXIT_PROJECT_BUDGET
BUILTIN_TOOLS = ["Read", "Write", "Edit", "Bash", "Glob", "Grep"]
# Everything else is denied by name. A bare-name deny rule removes the tool from
# the model's context entirely rather than denying it at call time, which is the
# behaviour we want: unavailable beats refused.
DENIED_TOOLS = ["WebSearch", "WebFetch", "NotebookEdit", "Task", "KillShell", "BashOutput"]
def _sdk() -> Any:
try:
import claude_agent_sdk # noqa: PLC0415
except ImportError as exc:
raise SystemExit(
"claude-agent-sdk is not installed.\n"
" pip install claude-agent-sdk\n"
"and authenticate with your subscription:\n"
" claude setup-token # then set CLAUDE_CODE_OAUTH_TOKEN"
) from exc
# Before any client is built: the SDK spawns the `claude` CLI without
# `creationflags`, and under the installed app (pythonw, no console) that
# is a black window titled "claude" over the workspace on every session.
from core import spawn # noqa: PLC0415
spawn.mask_sdk_console()
return claude_agent_sdk
def system_prompt() -> str:
"""The prompt, plus the two paths it cannot state for itself, plus memory.
`prompts/system.md` says to reach for "the skills in `skills/`", which is a
correct instruction only while the workspace and the installation are the
same folder. With the workspace pointed elsewhere the agent's cwd has no
`skills/` in it, and the model would go looking for a directory that is one
`..` away with nothing to say so.
Appended rather than templated into the file: the prompt is a document
someone edits, and a substitution marker in it is one more thing an edit can
break. The block is only added when the two directories actually differ, so
the standard install's prompt is byte-for-byte what the file says.
**The project's `MEMORY.md` is appended here, and here is why that is the
right seam.** Every session builds its options through `build_options`, and a
compaction builds a *fresh* client through the same function -- so putting
memory in the system prompt means the one thing a compaction must not lose is
the one thing it cannot lose, without `core/compaction.py` knowing anything
about projects. The handover note carries the conversation; this carries what
the project knows, and the second of those should not depend on a model
remembering to write it down in the first.
"""
text = paths.prompt_path().read_text(encoding="utf-8")
root, skills = paths.root(), paths.skills_dir()
appended: list[str] = []
if skills.parent != root:
appended.append(
"## Where things are\n\n"
f"Your working directory is the workspace: `{root}`. The ledger, notebooks, "
"notes and figures are there and every path you write should be relative to it.\n\n"
f"Grad itself is installed at `{paths.install_dir()}`, which is a separate "
f"folder you do not write to. The skills are there: `{skills}`."
)
memory = project_memory_block()
if memory:
appended.append(memory)
# Byte-for-byte when there is nothing to add, which is a property worth
# keeping rather than an accident: the prompt is a document someone edits and
# diffs, and a trailing newline this function invented would show up in every
# comparison against the file.
if not appended:
return text
return "\n\n".join([text.rstrip(), *appended]) + "\n"
def project_memory_block() -> str:
"""The selected project's memory, or "" -- and never an exception.
Wrapped this tightly because it runs before the first turn of every session
on both surfaces. A project directory that does not exist, a ledger that
cannot be read, a `MEMORY.md` locked by an editor: each of those is a session
that starts without memory, which is exactly the session everyone had before
this existed. None of them is a session that fails to start.
"""
try:
from core import budget, projects # noqa: PLC0415
cfg = config_mod.load()
limit = int(cfg.get("agent", "memory_max_chars", projects.MEMORY_MAX_CHARS))
return projects.prompt_block(budget.current_project(), max_chars=max(0, limit))
except Exception: # noqa: BLE001 - see the docstring
logging.getLogger("grad").debug("no project memory block", exc_info=True)
return ""
def build_options(
cfg: Any,
*,
permission_mode: str | None = None,
resume: str | None = None,
resume_at: str | None = None,
drops_turn: str | None = None,
) -> Any:
"""The options one session runs under.
`resume` is an SDK session id, and passing it is the difference between
reopening a transcript and reopening a *conversation*: without it the model
starts with no memory of the turns the window is showing above the composer,
which is a worse failure than not resuming at all because nothing on screen
says so. `ui/sessions.py` records the id and reports when it does not have
one.
`resume_at` is how a rewind reaches the model. It names the last transcript
entry to load, so the conversation comes back without the turns after it --
which is the only mechanism there is for this: the control protocol has no
"forget the last turn", and `core/rewind.py` explains why a rewind that only
cleaned the screen would be the worse half of the feature.
**Deliberately not `fork_session`.** Forking is the SDK's suggested
companion to `resume_at` and it remaps every uuid into a new session, which
would invalidate the anchor recorded on every turn above the rewind point --
so the *second* rewind in a session would find nothing to resume at and
silently degrade to a transcript-only one. Resuming in place keeps the ids
that `core/rewind.py:anchor_in` matches on, and the abandoned branch stays in
the SDK's own transcript either way.
`drops_turn` is validation, not behaviour: it names the user prompt the
truncation intends to discard and the CLI refuses the resume if anything
else would go with it -- a queued message, a wake the session absorbed
mid-turn. Optional, and a rewind that cannot identify one runs unvalidated
rather than not running.
"""
sdk = _sdk()
mode = permission_mode or str(cfg.get("agent", "permission_mode", "dontAsk"))
hook_matchers = {
"PreToolUse": [sdk.HookMatcher(matcher="Bash", hooks=[hooks.pre_tool_use])],
"Stop": [sdk.HookMatcher(hooks=[hooks.stop])],
}
options: dict[str, Any] = {
"resume": resume,
"model": cfg.model_for("research"),
"system_prompt": system_prompt(),
"allowed_tools": BUILTIN_TOOLS,
"disallowed_tools": DENIED_TOOLS,
"permission_mode": mode,
"cwd": str(paths.root()),
# The other half of `cwd`. Without it the agent's shell resolves `python`
# through the launcher's ambient PATH, which on this platform is not the
# environment Grad is running in -- see `interpreter_env`.
"env": interpreter_env(),
"hooks": hook_matchers,
# Off by default in the SDK, and the default is why an answer used to
# arrive in one lump: without it `receive_response` yields nothing until
# a whole `AssistantMessage` is finished. With it the same turn also
# emits `StreamEvent`s carrying token deltas. `TextStream` is what turns
# the two into one transcript -- see the warning in its docstring.
"include_partial_messages": True,
}
options.update(rewind_option(sdk, resume_at, drops_turn))
options.update(checkpointing_option(sdk))
options.update(thinking_option(cfg, sdk))
# How hard it thinks, from `core/effort.py`. Applied here rather than
# anywhere later because there is nowhere later: the SDK exposes no control
# request for it, so the level a client runs at is fixed when the client is
# built. `ui/app.py:Session.apply_effort` is what turns a change into a
# rebuild.
options.update(effort.option(cfg, sdk))
return sdk.ClaudeAgentOptions(**options)
def interpreter_env() -> dict[str, str]:
"""The environment the agent's shell runs in, so that `python` means this one.
`cwd` was the only thing this session told its shell about the world, and
`PATH` is the other half of that sentence. A `Bash` call inherits the
launcher's ambient environment, and `core/spawn.py:console_script` already
documents what that means on the platform this ships to: Grad is started
from a shortcut pointing at `.venv\\Scripts\\pythonw.exe`, Explorer hands it
the machine's environment, and so the interpreter is the venv's while `PATH`
is not. That reasoning was applied to `shutil.which` there and never to the
shell the model types into.
**What it costs is worse than a missing package.** Every tool in the system
prompt is spelled `python -m tools.<name>`, so the resolution of the bare
word `python` decides *which Grad* the agent's own instrument panel is. On a
machine carrying a second one on `PATH` -- a stale global install, an
editable checkout -- the answer is not this one, and the app then reads the
workspace through one installation while the agent writes it through
another. Nothing in either would ever name that disagreement.
So the interpreter running this process wins, and it wins by being *first*
rather than by being alone: `PATH` is prepended to, never replaced, because
the agent legitimately needs `git`, `docker` and `latexmk` and this is not
the place to decide what a research machine has on it.
Three details are load-bearing.
`sysconfig` is asked for the scripts directory rather than taking
`sys.executable`'s parent, because the two differ on a non-venv install:
`python.exe` sits in `C:\\Python314` and `pip.exe` in `C:\\Python314
\\Scripts`, and a fix that does not cover `pip` does not cover the command
people actually get wrong.
Nothing here is venv-specific, deliberately. The invariant is "the agent's
`python` is the `python` running Grad", which is the right one under a venv,
a conda environment or a bare system install -- and a check for a venv would
turn the third case into a silent no-op.
`PYTHONPATH` is set only when Grad is not an installed distribution, which
is the checkout someone runs with `python agent.py` and never pip-installed.
Unconditionally exporting the install directory would put `config`, `data`,
`notes` and `figures` on the import path as namespace packages, and `import
data` is a thing research code genuinely does.
`PYTHONUTF8` comes from `core/spawn.py:utf8_env`, which explains at length
why reading arXiv LaTeX on a Windows machine crashes twice over and why
remembering `encoding="utf-8"` is not enough to stop it. It is here rather
than in the prompt for the reason the gates are programs rather than
sentences: the model does get this right most of the time, and most of the
time is not a property.
"""
import sysconfig # noqa: PLC0415 - only this function needs it
from core import spawn # noqa: PLC0415
scripts = sysconfig.get_path("scripts") or str(Path(sys.executable).parent)
ambient = os.environ.get("PATH", "")
# Prepend rather than append, and skip the work when it is already in front:
# a session rebuilt by a compaction runs this again, and PATH should not grow
# a copy of the same directory once per compaction.
parts = ambient.split(os.pathsep) if ambient else []
if not parts or Path(parts[0] or ".") != Path(scripts):
ambient = os.pathsep.join([scripts, *parts]) if parts else scripts
env = {"PATH": ambient, **spawn.utf8_env()}
# `pip` and `uv` both read this, and an ambient one pointing at a *different*
# environment is the exact confusion this function exists to end -- so it is
# set to what is true here rather than left to whatever Explorer passed in.
if sys.prefix != sys.base_prefix:
env["VIRTUAL_ENV"] = sys.prefix
if not _installed_as_distribution():
install = str(paths.install_dir())
existing = os.environ.get("PYTHONPATH", "")
env["PYTHONPATH"] = os.pathsep.join([install, existing]) if existing else install
return env
def _installed_as_distribution() -> bool:
"""Is this Grad on the interpreter's path by installation rather than by cwd?
False is the safe answer on any failure: it adds `PYTHONPATH`, and a
redundant entry costs nothing next to an agent whose every tool call raises
`ModuleNotFoundError` because its cwd is the workspace and the workspace has
no `tools/` in it.
"""
try:
import importlib.metadata as md # noqa: PLC0415
md.distribution("grad")
return True
except Exception: # noqa: BLE001 - see the docstring
return False
def rewind_option(sdk: Any, resume_at: str | None, drops_turn: str | None) -> dict[str, Any]:
"""Where a resumed conversation should stop loading, when a rewind says so.
Feature-detected like `thinking_option`, and for the same reason: these two
options are newer than the rest of what this file passes, and an SDK without
them must give a session that ignores the rewind rather than one that cannot
be built at all. `Session.rewind_to` reads the transcript's half of the
rewind back off the return value, so the degradation is reported instead of
being left to be discovered when the model remembers a turn that is no longer
on screen.
`drops_turn` is dropped rather than sent alone when there is nothing to
anchor it to: it is a claim *about* a truncation, and on its own it would
describe one that is not happening.
"""
if not resume_at:
return {}
fields = _option_fields(sdk)
if "resume_session_at" not in fields:
return {}
options: dict[str, Any] = {"resume_session_at": resume_at}
if drops_turn and "resume_drops_turn" in fields:
options["resume_drops_turn"] = drops_turn
return options
def checkpointing_option(sdk: Any) -> dict[str, Any]:
"""Ask the CLI to keep a copy of a file before the agent changes it.
This is the third half of a rewind. `core/rewind.py` moves the transcript,
`resume_session_at` moves the model's memory, and until this flag neither of
them moved the *work*: a turn that rewrote a training script and was then
rewound left the rewritten script on disk, with a conversation that no longer
contained the instruction that produced it. The two surfaces the user can see
agreed with each other and disagreed with the filesystem, which is the worst
of the three arrangements.
Feature-detected, like `thinking_option` and `rewind_option` above, and for
the same reason: this option is newer than most of what this file passes, and
an SDK without it must give a session that cannot restore files rather than
one that cannot be built.
**What it does not cover.** The backups are taken by the CLI around its own
file-editing tools, so `Write` and `Edit` are the case it is for. Most of
what this agent does is a `Bash` command -- and a file a *command* wrote is
outside that, as is anything on the other side of a submitter. That is not a
gap worth closing here: the ledger is append-only precisely so the record of
a run cannot be rewound, and a rewind that reached into `ledger/runs.jsonl`
would be undoing evidence rather than work. `ui/app.py:rewind_to` words its
result on what actually moved rather than on what was asked for.
Deliberately not combined with `session_store`, which the SDK rejects
outright (`_internal/session_store_validation.py`). Nothing here sets one.
"""
if "enable_file_checkpointing" not in _option_fields(sdk):
return {}
return {"enable_file_checkpointing": True}
def checkpointing_supported(sdk: Any = None) -> bool:
"""Can this SDK put files back? Asked before anything promises it.
Separate from the option for the reason `rewind_supported` is separate from
`rewind_option`: the answer decides what the rewind's result *says*, and a
message claiming the work went back when it did not is the failure
`core/rewind.py` is written to make impossible.
"""
if sdk is None:
try:
sdk = _sdk()
except BaseException: # noqa: BLE001 - `_sdk` exits rather than raising ImportError
return False
return "enable_file_checkpointing" in _option_fields(sdk)
def _option_fields(sdk: Any) -> set[str]:
"""The option names this SDK accepts, or nothing if it cannot be asked."""
try:
return {field.name for field in dataclasses.fields(sdk.ClaudeAgentOptions)}
except TypeError:
return set()
def rewind_supported(sdk: Any = None) -> bool:
"""Can the installed SDK put a conversation back to an earlier point?
Asked separately from `rewind_option` because the answer is needed before
the option is built: an anchor the SDK will not accept is a rewind that
cleans the screen and leaves the model remembering everything, and telling
someone their memory went back when it did not is the exact failure
`core/rewind.py` is written to avoid. The caller words its result on this.
Feature-detected rather than pinned to an SDK floor, which is deliberate and
is the same call `thinking_option` makes. A hard minimum would refuse to
start for want of a capability the rest of the app does not need, on the
surface -- permission-mode names and option sets -- that has moved most
between releases. Degrading and saying so costs one sentence; a floor costs
everyone on an older SDK the whole application.
Answers False rather than raising when there is no SDK at all: the UI is
allowed to load without one, and a window asking whether a button should
promise something is not the place to discover it is missing.
"""
if sdk is None:
try:
sdk = _sdk()
except BaseException: # noqa: BLE001 - `_sdk` exits rather than raising ImportError
return False
return "resume_session_at" in _option_fields(sdk)
def thinking_option(cfg: Any, sdk: Any) -> dict[str, Any]:
"""Ask for the reasoning as text, when the installed SDK can be asked.
Capturing thinking blocks is not enough to *have* any: Opus 4.7+ defaults
`display` to "omitted" and sends them with a signature and no text. So the
chat window's reasoning switch had nothing to reveal no matter how correctly
the stream was read -- the bug was one flag away from the feature, and it
looked exactly like a toggle that did nothing.
Feature-detected rather than assumed, for the same reason `agent.py probe`
exists: this option is newer than the permission mode and the SDK's shape has
changed between releases. An SDK without it gets the options it understands
and a session with no reasoning, which is what it would have had anyway.
"""
display = str(cfg.get("agent", "reasoning", "summarized")).lower()
if display not in ("summarized", "omitted"):
display = "summarized"
try:
fields = {f.name for f in dataclasses.fields(sdk.ClaudeAgentOptions)}
except TypeError:
# Not a dataclass. `dataclasses.fields` raises rather than returning
# empty, and this whole function exists to tolerate an SDK whose options
# object is not the shape we expect -- so the release that changes it to
# an ordinary class or a TypedDict must degrade to "no reasoning
# settings", exactly as an SDK without the field already does, rather
# than take the app down before the first turn.
return {}
if "thinking" not in fields:
return {}
# `adaptive` rather than a fixed budget: the model decides how much thinking
# a turn is worth, which is the right call for a session that ranges from
# "what is in the ledger" to a campaign design.
return {"thinking": {"type": "adaptive", "display": display}}
def preflight_environment() -> dict[str, Any]:
"""Checks that must pass before the first turn.
ANTHROPIC_API_KEY outranks CLAUDE_CODE_OAUTH_TOKEN in the credential chain,
so a stray export silently bills the Developer Platform instead of the
subscription. It is removed here rather than warned about.
Then the complement, and the order between the two is the point: the scrub
takes out what must not be there, and `hydrate_environment` supplies the one
thing that must -- the subscription token, from the credential store, for
the app that was launched from a shortcut and never saw an `export`. Both
entry points reach this function (`run_session` and `ui/app.py`'s client
start), which is why the bridge belongs here rather than in either one.
"""
from core import budget # noqa: PLC0415
removed = credentials.scrub_environment()
# Read before hydrating, so the report can distinguish a token someone
# exported from one this just fetched. They authenticate identically; they
# are very different answers to "why is it using that account?".
ambient = bool(os.environ.get("CLAUDE_CODE_OAUTH_TOKEN"))
hydrated = credentials.hydrate_environment()
cfg = config_mod.load()
project_id = budget.current_project()
shell = interpreter_env()
return {
"removed_env": removed,
# What the agent's own `python` resolves to, which is the diagnostic that
# would have found the defect this reports on: `prompts/system.md` spells
# every tool `python -m tools.<name>`, so a second Grad earlier on PATH
# means the app and the agent are reading the same workspace through
# different installations. Names and booleans only -- this output goes
# into bug reports.
"interpreter": sys.executable,
"shell_path_head": shell["PATH"].split(os.pathsep)[0],
"venv": sys.prefix if sys.prefix != sys.base_prefix else None,
"installed_as_distribution": _installed_as_distribution(),
"oauth_token_present": bool(os.environ.get("CLAUDE_CODE_OAUTH_TOKEN")),
"oauth_token_source": (
"environment" if ambient else ("credential store" if hydrated else "absent")
),
"workspace": str(paths.root()),
"models": cfg.models(),
# Read from ledger/.current_project, not from the environment -- the
# scrub above is exactly why the selection is a file (§15).
"project": project_id,
"project_status": budget.status(project_id) if budget.exists(project_id) else None,
"note": (
"auth should be subscription-backed; confirm with `claude /status`. "
"--bare mode does not read CLAUDE_CODE_OAUTH_TOKEN, so this runs non-bare."
),
}
# ---------------------------------------------------------------------------
# session
# ---------------------------------------------------------------------------
async def run_session(prompt: str | None, *, once: bool) -> int:
sdk = _sdk()
cfg = config_mod.load()
paths.ensure_workspace()
env = preflight_environment()
if env["removed_env"]:
print(f"[grad] removed from the environment: {', '.join(env['removed_env'])}", file=sys.stderr)
# Held in a variable rather than in an `async with`, because compacting
# replaces it: the note is written by the outgoing session and the fresh one
# is built to hold it. A context manager binds the name for the whole block
# and there would be no way to swap what it holds -- which is how the CLI
# would have ended up as the surface that cannot compact, and the two
# surfaces disagreeing about a rule is the failure `drive_turn`'s docstring
# is about.
client = await _connect(sdk, cfg)
#: The handover note from a compaction, waiting for the next prompt to ride
#: in front of. Sending it as a turn of its own would spend a round-trip to
#: produce an answer nobody asked for.
seed: str | None = None
try:
if prompt:
ran, seed = await _turn(client, prompt, seed)
if once:
# No compaction on a one-shot: the session ends here, so the only
# thing a compaction could buy is a summary nothing will read.
return 0 if ran else EXIT_PROJECT_BUDGET
client, seed = await _maybe_compact(sdk, cfg, client, seed)
while True:
try:
# In a worker thread: a bare input() blocks the event loop, and
# the SDK client cannot service its transport while it waits --
# so streaming, keepalives, and interrupts stall for the whole
# idle period between turns.
line = (await asyncio.to_thread(input, "\n> ")).strip()
except (EOFError, KeyboardInterrupt):
print()
return 0
if not line:
continue
if line in ("exit", "quit"):
return 0
_, seed = await _turn(client, line, seed)
client, seed = await _maybe_compact(sdk, cfg, client, seed)
finally:
await _disconnect(client)
async def _connect(sdk: Any, cfg: Any, *, resume: str | None = None) -> Any:
client = sdk.ClaudeSDKClient(options=build_options(cfg, resume=resume))
await client.__aenter__()
return client
async def _disconnect(client: Any) -> None:
"""Exit a client's context. Never raises on the way out."""
if client is None:
return
try:
await client.__aexit__(None, None, None)
except Exception: # noqa: BLE001 - shutdown must not raise
pass
async def _maybe_compact(sdk: Any, cfg: Any, client: Any, seed: str | None) -> tuple[Any, str | None]:
"""Compact between turns when the context has passed the threshold.
The CLI's half of what `ui/app.py:Session.maybe_compact` does, and the same
order for the same reason: the note is written while the outgoing session
still remembers everything, and only then is the client replaced.
A failure here returns the client unchanged. An oversized conversation is a
cost; a session taken down between turns by its own housekeeping is a loss.
"""
from core import compaction # noqa: PLC0415
if not compaction.threshold(cfg):
return client, seed
reader = getattr(client, "get_context_usage", None)
if reader is None:
return client, seed
try:
usage = await reader()
except Exception: # noqa: BLE001 - no reading is not a reason to compact
return client, seed
if not compaction.should_compact(usage, cfg):
return client, seed
before = compaction.context_tokens(usage)
print(f"\n[grad] compacting at {before:,} tokens…", file=sys.stderr)
try:
handoff = await compaction.write_handoff(client, drive_turn)
except BudgetRefused as exc:
# No carve-out: a compaction is a model call and the allocation applies.
print(f"[grad] cannot compact: {exc.refusal['message']}", file=sys.stderr)
return client, seed
except Exception as exc: # noqa: BLE001 - the conversation survives a failed compaction
print(f"[grad] could not compact ({type(exc).__name__}); carrying on", file=sys.stderr)
return client, seed
await _disconnect(client)
# `resume` is deliberately not passed. Resuming would restore the very
# conversation this just summarised, making the whole operation a cost with
# no effect.
fresh = await _connect(sdk, cfg)
print("[grad] compacted — the agent now knows this session by its handover note", file=sys.stderr)
return fresh, compaction.seed_message(handoff["note"], tokens_before=before)
def check_turn_budget() -> dict[str, Any] | None:
"""Refuse the *next* turn when the project is out of token allocation.
HANDOFF-2 §15, and the honesty is the point: tokens are consumed
continuously inside a turn and there is no way to refuse mid-turn, so
**token budgets are enforced to a granularity of one turn's overrun.** This
check is our code end to end -- it depends on no SDK behaviour -- and it runs
before `query`, not after.
Returns a refusal payload, or None to proceed.
"""
try:
from core import budget # noqa: PLC0415
project_id = budget.current_project()
if not project_id or not budget.exists(project_id):
return None
state = budget.status(project_id)
except Exception as exc: # noqa: BLE001 - accounting must never strand a session
# Fails open, and says so. This is the *only* mechanism that bounds
# token spend before a turn; if it cannot read the ledger, the honest
# report is that the turn is going out ungated.
print(
f"[grad] token budget check failed ({type(exc).__name__}: {exc}); "
"this turn is not gated",
file=sys.stderr,
)
return None
tokens = state["resources"]["quota_tokens"]
if not tokens["over"]:
return None
overrun = tokens["spent"] - float(tokens["ceiling"])
return {
"project": project_id,
"resource": "quota_tokens",
"spent": tokens["spent"],
"ceiling": tokens["ceiling"],
"overrun": overrun,
"message": (
f"project {project_id} has used {tokens['spent']:,} of its "
f"{int(tokens['ceiling']):,} token allocation -- {overrun:,.0f} over. "
"Refusing the next turn; the turn that crossed the ceiling was allowed to "
"finish, because there is no way to refuse mid-turn."
),
"fix": (
f"python -m tools.budget raise --project {project_id} "
"--quota-tokens <new ceiling> --json"
),
}
class BudgetRefused(Exception):
"""Raised by `drive_turn` when the project is out of token allocation.
Carries the payload so a caller can render it: the CLI prints it, the UI
puts it in the transcript.
"""
def __init__(self, refusal: dict[str, Any]) -> None:
super().__init__(refusal["message"])
self.refusal = refusal
async def drive_turn(
client: Any,
prompt: str,
stream: Any,
*,
on_chunk: Any = None,
on_session_id: Any = None,
on_uuid: Any = None,
session: str | None = None,
stage: str = quota_log.STAGE_MAIN,
role: str = "research",
) -> dict[str, Any]:
"""One turn, for every surface that runs one.
The CLI loop and the UI's `Session.ask` were the same loop written twice,
and only one of them checked the budget or recorded what the turn spent --
so everything done through the desktop app, which is the primary surface,
accrued no tokens in `ledger/quota.jsonl` and passed no ceiling. The README
said the allocation is checked "before issuing the next turn"; that was true
of `python agent.py` and false of `python agent.py --ui`. One driver, so
there is one answer.
`on_chunk` is called with each newly-visible piece of text; the UI passes
nothing because its renderer reads `stream.blocks` on a timer instead.
`on_session_id` is called the moment the SDK names this conversation, and it
exists because the return value is not reached on every path this function
can take. A turn that is *interrupted* raises, so a caller reading the id
off the return value learned it only for turns that finished -- and an
interrupted turn is precisely the one after which the client is rebuilt, so
that was the case where losing the id cost the whole conversation.
`on_uuid` is called with each transcript entry the SDK names, and it is a
callback for the same reason and with a sharper version of the same
consequence. The *last* one of a turn is what `resume_session_at` takes to
rewind to the end of that turn (`core/rewind.py`), and a turn that failed
half-way is the single most likely thing anyone will want to rewind past --
so reading it off a return value that a failed turn never reaches would lose
the anchor on precisely the turns the feature exists for.
`stage` and `role` decide where the turn's tokens land in `ledger/quota.jsonl`.
They default to the conversation, and the one caller that overrides them is
`core/compaction.py`: a compaction is a model call this system makes on its
own initiative, and folding its cost into `main` would hide precisely the
number that says whether the threshold is set correctly.
"""
refusal = check_turn_budget()
if refusal:
raise BudgetRefused(refusal)
await client.query(prompt)
sdk_session_id: str | None = None
last_usage: Any = None
recorded = None
try:
async for message in client.receive_response():
# Whatever has not been printed yet -- a token as it arrives, the
# tail of a message that was never streamed, or a line naming a tool
# call. Never both halves of the same text.
chunk = stream.feed(message)
if chunk and on_chunk is not None:
on_chunk(chunk)
# Captured from the stream rather than asked for: the SDK assigns
# it, and this is the id `resume` takes when a session is reopened.
# A resumed conversation can be given a new id, so the latest wins.
candidate = getattr(message, "session_id", None)
if isinstance(candidate, str) and candidate:
if candidate != sdk_session_id and on_session_id is not None:
on_session_id(candidate)
sdk_session_id = candidate
# The id of this entry in the SDK's own transcript, which is what a
# rewind resumes at. Reported as it arrives and never compared
# against what came before: the last one wins because the anchor
# wanted is the *end* of the turn, whatever kind of message that
# turned out to be. `ResultMessage` is usually it and carries one.
entry = getattr(message, "uuid", None)
if isinstance(entry, str) and entry and on_uuid is not None:
on_uuid(entry)
# The *last* usage seen, recorded once after the loop -- not one
# record per message. `ResultMessage` arrives last and carries the
# turn's cumulative usage, so summing every message that has a
# `usage` attribute would count the same tokens twice.
usage = getattr(message, "usage", None)
if usage is not None:
last_usage = usage
finally:
# In a `finally` because a turn that died half-way still spent what it
# spent. Letting the exception skip this would make a failing session
# the cheapest way to run untracked -- the accounting would be missing
# exactly the turns most worth accounting for.
if last_usage is not None:
recorded = quota_log.from_sdk_usage(
stage, last_usage, model=None, role=role, session=session
)
return {"sdk_session_id": sdk_session_id, "quota": recorded}
async def _turn(client: Any, prompt: str, seed: str | None = None) -> tuple[bool, str | None]:
"""Run one turn. Returns whether it ran, and the seed still owed.
`seed` is a handover note from a compaction, prepended to this prompt rather
than sent as a turn of its own. It is returned unconsumed when the turn does
not run, because a note dropped by a refused turn is the whole memory of
everything the compaction discarded.
"""
stream = TurnStream()
sent = f"{seed}\n\n---\n\n{prompt}" if seed else prompt
try:
await drive_turn(
client, sent, stream, on_chunk=lambda c: print(c, end="", flush=True)
)
except BudgetRefused as exc:
print(
f"\n[grad] {exc.refusal['message']}\n[grad] fix: {exc.refusal['fix']}",
file=sys.stderr,
)
return False, seed
print()
return True, None
def _text_of(message: Any) -> str:
content = getattr(message, "content", None)
if isinstance(content, str):
return content
if isinstance(content, list):
return "".join(getattr(b, "text", "") or "" for b in content)
return ""
def _delta_of(message: Any) -> str:
"""The visible text a partial-message stream event carries, if any.
`StreamEvent.event` is the raw Anthropic streaming event, so this is a
filter as much as an accessor: only `content_block_delta` carrying a
`text_delta` is answer text. Thinking deltas and tool-input deltas are
excluded deliberately, because `_text_of` excludes their finished blocks too
-- a `ThinkingBlock` has `.thinking`, not `.text`. Letting them through here
would make the stream say something the settled message does not.
"""
event = getattr(message, "event", None)
if not isinstance(event, dict) or event.get("type") != "content_block_delta":
return ""
delta = event.get("delta")
if not isinstance(delta, dict) or delta.get("type") != "text_delta":
return ""
text = delta.get("text")
return text if isinstance(text, str) else ""
def _thinking_delta_of(message: Any) -> str:
"""The reasoning a partial-message stream event carries, if any.
The mirror of `_delta_of`, and a separate function rather than a parameter on
it: the two feed different halves of the transcript, and the whole reason
`_delta_of` filters as hard as it does is that mixing them makes the stream
say something the settled message does not.
"""
event = getattr(message, "event", None)
if not isinstance(event, dict) or event.get("type") != "content_block_delta":
return ""
delta = event.get("delta")
if not isinstance(delta, dict) or delta.get("type") != "thinking_delta":
return ""
text = delta.get("thinking")
return text if isinstance(text, str) else ""
def _thinking_of(message: Any) -> str:
"""The reasoning in a finished message.
A `ThinkingBlock` carries `.thinking`, not `.text`, which is exactly why
`_text_of` misses it -- and why the reasoning needed a second pair of
accessors rather than a looser filter on the first.
"""
content = getattr(message, "content", None)
if not isinstance(content, list):
return ""
return "".join(getattr(b, "thinking", "") or "" for b in content)
class TextStream:
"""One turn's visible text, assembled from deltas *and* finished messages.
`include_partial_messages` makes the SDK emit both halves of the same text:
a run of `text_delta` events, and then the `AssistantMessage` that contains
all of it. **Appending both is the bug this class exists to prevent** -- it
is the obvious way to write the loop, and it makes every answer appear
twice.
So a finished message *replaces* the deltas that built it rather than
following them. That ordering also makes the finished message authoritative:
if the two ever disagree -- a dropped event, a turn resumed from cache, a
message the SDK never streamed -- what stays on screen is the message, not
the reconstruction. A turn is many messages, so this repeats per message,
which is why `_streamed` is reset each time rather than once at the end.
`feed` returns only the text that has not been shown yet, so a CLI can print
its return value directly; `text` is the whole answer so far, for a UI that
re-renders from it.
The two accessors are parameters because the *reasoning* half of a turn
arrives the same way and has the same trap: `thinking_delta` events followed
by a `ThinkingBlock` containing all of them. One class, given the other pair
of accessors, is what keeps the no-duplication rule stated once.
"""
def __init__(self, delta_of: Any = None, whole_of: Any = None) -> None:
self.text = ""
#: The tail of `text` contributed by deltas since the last finished
#: message -- the part a finished message is entitled to overwrite.
self._streamed = ""
self._delta_of = delta_of or _delta_of
self._whole_of = whole_of or _text_of
def feed(self, message: Any) -> str:
delta = self._delta_of(message)
if delta:
self.text += delta
self._streamed += delta
return delta
text = self._whole_of(message)
# A message with no text at all -- a tool result, a system message, the
# final result -- must leave a half-streamed block alone.
if not text:
return ""
if text.startswith(self._streamed):
unseen = text[len(self._streamed) :]
self.text += unseen
else:
self.text = self.text[: len(self.text) - len(self._streamed)] + text
unseen = ""
self._streamed = ""
return unseen
# ---------------------------------------------------------------------------
# tool calls
# ---------------------------------------------------------------------------
#: What one call contributes to a transcript, at most. A `Read` of a long file
#: or a training log is tens of thousands of characters, and every one of them
#: would be held for the life of the session, written to the transcript file on
#: settle, and drawn again on restore. A card is a record that the call happened
#: and how it went, not a second copy of its output.
RESULT_CHARS = 2000
RESULT_LINES = 40
#: Which input key says what a call was *on*. Anything not listed here falls
#: back to `SUBJECT_KEYS`, then to the first short string in the input -- an
#: `Edit` carries its whole replacement text, and a card head that is a wall of
#: source is worse than one that is empty.
TOOL_SUBJECT = {
"Bash": "command",
"Read": "file_path",
"Write": "file_path",
"Edit": "file_path",
"Glob": "pattern",
"Grep": "pattern",
}
SUBJECT_KEYS = ("command", "file_path", "path", "pattern", "query", "url", "prompt")
#: Per-row limits for the rest of a call's input. Six rows of two lines is a
#: card you can read at a glance; the whole input is not.
ROW_CHARS = 200
ROW_LINES = 2
MAX_ROWS = 6
def clip(text: str, *, chars: int = RESULT_CHARS, lines: int = RESULT_LINES) -> str:
"""`text`, bounded -- and saying what it dropped rather than trailing off.
ASCII on purpose: this can reach a Windows console, where a stray `…` is a
`UnicodeEncodeError` that would take the turn down.
"""
if not text:
return ""
split = text.splitlines()
dropped_lines = max(0, len(split) - lines)
out = "\n".join(split[:lines])
dropped_chars = max(0, len(out) - chars)
out = out[:chars]
if dropped_chars:
out += f"\n... +{dropped_chars:,} more characters"
if dropped_lines:
out += f"\n... +{dropped_lines:,} more lines"
return out
def _one_line(text: str, limit: int = 120) -> str:
"""A subject collapsed onto one line, for a card head."""
flattened = " ".join(str(text).split())
return flattened if len(flattened) <= limit else flattened[: limit - 3] + "..."
def describe_tool(name: str, tool_input: dict[str, Any]) -> tuple[str, str]:
"""`(subject_key, subject)` -- what this call was on, and under which key."""
keys = (TOOL_SUBJECT[name],) if name in TOOL_SUBJECT else SUBJECT_KEYS
for key in keys:
value = tool_input.get(key)
if isinstance(value, str) and value.strip():
return key, value
for key, value in tool_input.items():
if isinstance(value, str) and value.strip() and len(value) <= ROW_CHARS:
return key, value
return "", ""
def _content_blocks(message: Any) -> list[Any]:
content = getattr(message, "content", None)
return content if isinstance(content, list) else []
def _tool_uses(message: Any) -> list[dict[str, Any]]:
"""Every `ToolUseBlock` in a finished message, as plain data.
Duck-typed rather than isinstance-checked so `ServerToolUseBlock` -- a tool
the API runs on the model's behalf -- draws the same card, and so this file
keeps working if the SDK renames a class.
"""
uses: list[dict[str, Any]] = []
for block in _content_blocks(message):
name = getattr(block, "name", None)