1212from dataclasses import dataclass , field
1313from datetime import datetime , timezone
1414from pathlib import Path
15- from typing import Any , AsyncIterator , Callable , Iterator , Optional , Protocol , Sequence
15+ from typing import Any , AsyncIterator , Callable , Iterator , Optional , Sequence
1616
1717import psutil
1818import yaml
4040_USER_PROMPT_FILENAME = "user_prompt.md"
4141_PROGRESS_ENV = "DSTACK_PRESET_PROGRESS_LOG"
4242_REDACTION = "[redacted]"
43+ # Replacing shorter values such as "1" or "false" corrupts unrelated diagnostics.
44+ _MIN_REDACTED_SUBSTRING_LENGTH = 8
4345_CLAUDE_TOOLS = "Bash,Read,Write,Edit,WebFetch,WebSearch,StructuredOutput"
4446_CLAUDE_EFFORT_LEVELS = ("low" , "medium" , "high" , "xhigh" , "max" )
4547_RESUME_DELAYS_SECONDS : tuple [int , ...] = (30 , 60 , 120 )
@@ -137,7 +139,6 @@ def agent_stderr_path(self) -> Path:
137139@dataclass
138140class PresetAgentSession :
139141 path : Path
140- timestamp : str
141142 debug : bool
142143 preset_id : str = ""
143144 # Whether progress lines echo to this process's console (a live attach), on
@@ -328,7 +329,6 @@ def create_preset_agent_session(
328329) -> PresetAgentSession :
329330 if configuration .name is None :
330331 raise CLIError ("The service name is required to save agent output" )
331- timestamp = datetime .now (timezone .utc ).strftime ("%Y%m%d-%H%M%S-%fZ" )
332332 parent = get_presets_dir ()
333333 path : Optional [Path ] = None
334334 try :
@@ -371,7 +371,7 @@ def create_preset_agent_session(
371371 shutil .rmtree (path , ignore_errors = True )
372372 raise CLIError (f"Could not create agent output under { parent } : { e } " ) from e
373373 assert path is not None
374- return PresetAgentSession (path = path , timestamp = timestamp , debug = debug , preset_id = preset_id )
374+ return PresetAgentSession (path = path , debug = debug , preset_id = preset_id )
375375
376376
377377def _get_claude_version (auth : "ClaudeAuth" ) -> Optional [str ]:
@@ -407,7 +407,7 @@ def _get_claude_auth_status(auth: "ClaudeAuth") -> dict[str, Any]:
407407
408408def load_resumable_agent_session (preset_id : str ) -> PresetAgentSession :
409409 path = get_presets_dir () / preset_id
410- session = PresetAgentSession (path = path , timestamp = "" , debug = False , preset_id = preset_id )
410+ session = PresetAgentSession (path = path , debug = False , preset_id = preset_id )
411411 manifest = session .read_manifest ()
412412 if not path .is_dir () or not manifest :
413413 raise CLIError (f"Unknown preset: { preset_id } " )
@@ -424,7 +424,6 @@ def load_resumable_agent_session(preset_id: str) -> PresetAgentSession:
424424 if not manifest .get ("claude_session_id" ):
425425 raise CLIError (f"Preset { preset_id } creation stopped before it started; create a new one" )
426426 session .debug = bool (manifest .get ("debug" ))
427- session .timestamp = str (manifest .get ("created_at" ) or "" )
428427 return session
429428
430429
@@ -455,7 +454,7 @@ def session_process_alive(manifest: dict[str, Any]) -> bool:
455454
456455def load_attachable_agent_session (preset_id : str ) -> PresetAgentSession :
457456 path = get_presets_dir () / preset_id
458- session = PresetAgentSession (path = path , timestamp = "" , debug = False , preset_id = preset_id )
457+ session = PresetAgentSession (path = path , debug = False , preset_id = preset_id )
459458 manifest = session .read_manifest ()
460459 if not path .is_dir () or not manifest :
461460 raise CLIError (f"Unknown preset: { preset_id } " )
@@ -481,14 +480,13 @@ def load_attachable_agent_session(preset_id: str) -> PresetAgentSession:
481480 f" stop or detach it there with Ctrl+C"
482481 )
483482 session .debug = bool (manifest .get ("debug" ))
484- session .timestamp = str (manifest .get ("created_at" ) or "" )
485483 return session
486484
487485
488486def load_agent_session (preset_id : str ) -> PresetAgentSession :
489487 """Loads a session of any status for read-only inspection (its log)."""
490488 path = get_presets_dir () / preset_id
491- session = PresetAgentSession (path = path , timestamp = "" , debug = False , preset_id = preset_id )
489+ session = PresetAgentSession (path = path , debug = False , preset_id = preset_id )
492490 if not path .is_dir () or not session .read_manifest ():
493491 raise CLIError (f"Unknown preset: { preset_id } " )
494492 return session
@@ -601,7 +599,7 @@ def iter_agent_sessions() -> Iterator[PresetAgentSession]:
601599 return
602600 for path in sorted (root .iterdir ()):
603601 if path .is_dir () and not path .name .startswith (("." , "models--" )):
604- yield PresetAgentSession (path = path , timestamp = "" , debug = False , preset_id = path .name )
602+ yield PresetAgentSession (path = path , debug = False , preset_id = path .name )
605603
606604
607605def find_session_name_claims (name : str ) -> list [PresetAgentSession ]:
@@ -882,7 +880,8 @@ def get_redacted_values(values: Sequence[str]) -> tuple[str, ...]:
882880def contains_redacted_value (value : Any , redacted_values : Sequence [str ]) -> bool :
883881 if isinstance (value , str ):
884882 return any (
885- value == redacted or (len (redacted ) >= 8 and redacted in value )
883+ value == redacted
884+ or (len (redacted ) >= _MIN_REDACTED_SUBSTRING_LENGTH and redacted in value )
886885 for redacted in redacted_values
887886 )
888887 if isinstance (value , dict ):
@@ -900,12 +899,25 @@ def redact(value: str, redacted_values: Sequence[str]) -> str:
900899 for redacted_value in redacted_values :
901900 if value == redacted_value :
902901 return _REDACTION
903- # Replacing short values such as "1" or "false" corrupts unrelated diagnostics.
904- if len (redacted_value ) >= 8 :
902+ if len (redacted_value ) >= _MIN_REDACTED_SUBSTRING_LENGTH :
905903 value = value .replace (redacted_value , _REDACTION )
906904 return value
907905
908906
907+ def redact_structure (value : Any , redacted_values : Sequence [str ]) -> Any :
908+ """Recursively redacts every string (including dict keys) in a JSON-like value."""
909+ if isinstance (value , str ):
910+ return redact (value , redacted_values )
911+ if isinstance (value , list ):
912+ return [redact_structure (item , redacted_values ) for item in value ]
913+ if isinstance (value , dict ):
914+ return {
915+ redact (key , redacted_values ): redact_structure (item , redacted_values )
916+ for key , item in value .items ()
917+ }
918+ return value
919+
920+
909921def _validate_control_socket_path (build_root : Path ) -> None :
910922 if IS_WINDOWS :
911923 return
@@ -1080,9 +1092,31 @@ def _prepare_subprocess_command(command: list[str]) -> list[str]:
10801092
10811093
10821094def _write_private_text (path : Path , content : str ) -> None :
1083- path .write_text (content , encoding = "utf-8" )
1084- if not IS_WINDOWS :
1085- path .chmod (0o600 )
1095+ # Atomic tmp + fsync + replace (mkstemp already creates the file 0600), so
1096+ # a crash mid-write cannot leave a truncated manifest or offsets file.
1097+ fd , temporary = tempfile .mkstemp (dir = path .parent , prefix = f".{ path .name } ." , suffix = ".tmp" )
1098+ try :
1099+ with os .fdopen (fd , "w" , encoding = "utf-8" ) as f :
1100+ f .write (content )
1101+ f .flush ()
1102+ os .fsync (f .fileno ())
1103+ try :
1104+ os .replace (temporary , path )
1105+ except PermissionError :
1106+ if not IS_WINDOWS :
1107+ raise
1108+ # A concurrent reader (a viewer polling the manifest) can hold the
1109+ # destination open without FILE_SHARE_DELETE; retry briefly, then
1110+ # prefer an in-place write over crashing the owner.
1111+ for _ in range (3 ):
1112+ time .sleep (0.01 )
1113+ with suppress (PermissionError ):
1114+ os .replace (temporary , path )
1115+ return
1116+ path .write_text (content , encoding = "utf-8" )
1117+ finally :
1118+ with suppress (FileNotFoundError ):
1119+ os .unlink (temporary )
10861120
10871121
10881122def _write_debug_trace (
@@ -1096,7 +1130,7 @@ def _write_debug_trace(
10961130 datetime .now (timezone .utc ).isoformat (timespec = "microseconds" ).replace ("+00:00" , "Z" )
10971131 )
10981132 try :
1099- event = _redact_trace_value (json .loads (text ), redacted_values )
1133+ event = redact_structure (json .loads (text ), redacted_values )
11001134 record = {"timestamp" : timestamp , "stream" : stream_name , "event" : event }
11011135 except json .JSONDecodeError :
11021136 record = {
@@ -1110,19 +1144,6 @@ def _write_debug_trace(
11101144 f .flush ()
11111145
11121146
1113- def _redact_trace_value (value : Any , redacted_values : Sequence [str ]) -> Any :
1114- if isinstance (value , str ):
1115- return redact (value , redacted_values )
1116- if isinstance (value , list ):
1117- return [_redact_trace_value (item , redacted_values ) for item in value ]
1118- if isinstance (value , dict ):
1119- return {
1120- redact (key , redacted_values ): _redact_trace_value (item , redacted_values )
1121- for key , item in value .items ()
1122- }
1123- return value
1124-
1125-
11261147@asynccontextmanager
11271148async def _session_tailers (
11281149 * ,
@@ -1182,7 +1203,7 @@ async def _collect_agent_output(
11821203 """Parses the agent's stream files until it exits; safe alongside a live
11831204 process or over the remains of a finished one."""
11841205 offset_store = _OffsetStore (agent_session .path / ".offsets.json" )
1185- stdout_output , stderr_output = await asyncio .gather (
1206+ stdout_output , _ = await asyncio .gather (
11861207 _read_process_stream (
11871208 stream = _FileLineReader (
11881209 workspace .agent_stdout_path ,
@@ -1208,12 +1229,9 @@ async def _collect_agent_output(
12081229 agent_session = agent_session ,
12091230 ),
12101231 )
1211- output = stdout_output
1212- if output .report_data is None :
1213- output .report_data = stderr_output .report_data
1214- if output .error is None :
1215- output .error = stderr_output .error
1216- return output
1232+ # stderr is tailed with parse_result=False — it feeds the debug trace and
1233+ # advances the persisted offset, but can never contribute report data.
1234+ return stdout_output
12171235
12181236
12191237async def attach_preset_agent (
@@ -1242,7 +1260,7 @@ def agent_alive() -> bool:
12421260
12431261async def _read_process_stream (
12441262 * ,
1245- stream : "_LineStream " ,
1263+ stream : "_FileLineReader " ,
12461264 stream_name : str ,
12471265 parse_result : bool ,
12481266 redacted_values : Sequence [str ],
@@ -1370,10 +1388,6 @@ def _terminate_windows_process_tree(pid: int) -> None:
13701388 psutil .wait_procs (alive , timeout = 3 )
13711389
13721390
1373- class _LineStream (Protocol ):
1374- async def readline (self ) -> bytes : ...
1375-
1376-
13771391class _FileLineReader :
13781392 """`readline()` over a growing file, so stream parsing survives CLI
13791393 restarts: offsets persist, and a later attach continues exactly where the
0 commit comments