Leon4gr45 commited on
Commit
f086669
·
verified ·
1 Parent(s): cfeb3a1

Upload folder using huggingface_hub (part 4)

Browse files
.gitattributes CHANGED
@@ -448,3 +448,20 @@ helpers/mcp_server.py.dox.md filter=lfs diff=lfs merge=lfs -text
448
  helpers/media_artifacts.py filter=lfs diff=lfs merge=lfs -text
449
  helpers/media_artifacts.py.dox.md filter=lfs diff=lfs merge=lfs -text
450
  helpers/message_queue.py filter=lfs diff=lfs merge=lfs -text
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
448
  helpers/media_artifacts.py filter=lfs diff=lfs merge=lfs -text
449
  helpers/media_artifacts.py.dox.md filter=lfs diff=lfs merge=lfs -text
450
  helpers/message_queue.py filter=lfs diff=lfs merge=lfs -text
451
+ helpers/message_queue.py.dox.md filter=lfs diff=lfs merge=lfs -text
452
+ helpers/messages.py filter=lfs diff=lfs merge=lfs -text
453
+ helpers/messages.py.dox.md filter=lfs diff=lfs merge=lfs -text
454
+ helpers/microsoft_tunnel.py filter=lfs diff=lfs merge=lfs -text
455
+ helpers/microsoft_tunnel.py.dox.md filter=lfs diff=lfs merge=lfs -text
456
+ helpers/migration.py filter=lfs diff=lfs merge=lfs -text
457
+ helpers/migration.py.dox.md filter=lfs diff=lfs merge=lfs -text
458
+ helpers/modules.py filter=lfs diff=lfs merge=lfs -text
459
+ helpers/modules.py.dox.md filter=lfs diff=lfs merge=lfs -text
460
+ helpers/network.py filter=lfs diff=lfs merge=lfs -text
461
+ helpers/network.py.dox.md filter=lfs diff=lfs merge=lfs -text
462
+ helpers/notification.py filter=lfs diff=lfs merge=lfs -text
463
+ helpers/notification.py.dox.md filter=lfs diff=lfs merge=lfs -text
464
+ helpers/parallel_tools.py filter=lfs diff=lfs merge=lfs -text
465
+ helpers/parallel_tools.py.dox.md filter=lfs diff=lfs merge=lfs -text
466
+ helpers/performance.py filter=lfs diff=lfs merge=lfs -text
467
+ helpers/performance.py.dox.md filter=lfs diff=lfs merge=lfs -text
helpers/message_queue.py.dox.md CHANGED
@@ -1,53 +1,3 @@
1
- # message_queue.py DOX
2
-
3
- ## Purpose
4
-
5
- - Own the `message_queue.py` helper module.
6
- - This module stores and drains queued user messages for a context.
7
- - Keep this file-level DOX profile synchronized with `message_queue.py` because this directory is intentionally flat.
8
-
9
- ## Ownership
10
-
11
- - `message_queue.py` owns the runtime implementation.
12
- - `message_queue.py.dox.md` owns durable notes about responsibilities, contracts, side effects, and verification for that implementation.
13
- - Top-level functions:
14
- - `get_queue(context: 'AgentContext') -> list`: Get current queue from context.data.
15
- - `_get_next_seq(context: 'AgentContext') -> int`: Get next sequence number.
16
- - `_sync_output(context: 'AgentContext')`: Sync queue to output_data for frontend polling.
17
- - `add(context: 'AgentContext', text: str, attachments: list[str] | None=..., item_id: str | None=...) -> dict`: Add message to queue. Attachments should be filenames, will be converted to full paths.
18
- - `remove(context: 'AgentContext', item_id: str | None=...) -> int`: Remove item(s). If item_id is None, clears all. Returns remaining count.
19
- - `pop_first(context: 'AgentContext') -> dict | None`: Remove and return first item.
20
- - `pop_item(context: 'AgentContext', item_id: str) -> dict | None`: Remove and return specific item.
21
- - `has_queue(context: 'AgentContext') -> bool`: Check if queue has items.
22
- - `log_user_message(context: 'AgentContext', message: str, attachment_paths: list[str], message_id: str | None=..., source: str=...)`: Log user message to console and UI. Used by message API and queue processing.
23
- - `send_message(context: 'AgentContext', item: dict, source: str=...)`: Send a single queued message (log + communicate).
24
- - `send_next(context: 'AgentContext') -> bool`: Send next queued message. Returns True if sent, False if queue empty.
25
- - `send_all_aggregated(context: 'AgentContext') -> int`: Aggregate and send all queued messages as one. Returns count of items sent.
26
- - Notable constants/configuration names: `QUEUE_KEY`, `QUEUE_SEQ_KEY`, `UPLOAD_FOLDER`.
27
-
28
- ## Runtime Contracts
29
-
30
- - Helper modules own reusable framework APIs and must preserve public callers unless all callers, tests, and docs are updated together.
31
- - Update this file whenever public functions, classes, persistence behavior, path/security assumptions, side effects, or cross-module contracts change.
32
- - Observed side-effect areas: filesystem reads, filesystem deletion, settings/state persistence.
33
- - Imported dependency areas include: `helpers`, `helpers.print_style`, `os`, `typing`, `uuid`.
34
-
35
- ## Key Concepts
36
-
37
- - Important called helpers/classes observed in the source: `context.set_data`, `get_queue`, `context.set_output_data`, `_sync_output`, `queue.pop`, `context.log.log`, `log_user_message`, `context.communicate`, `pop_first`, `has_queue`, `join`, `context.get_data`, `att.startswith`, `_get_next_seq`, `uuid.uuid4`, `UserMessage`, `send_message`, `guids.generate_id`, `os.path.basename`, `PrintStyle`.
38
- - Keep request/response, tool, or helper semantics documented here at the same time as source changes.
39
-
40
- ## Work Guidance
41
-
42
- - Preserve public helper APIs used by core code and plugins unless every caller is updated.
43
- - Keep path, auth, secret, persistence, network, and subprocess behavior explicit and bounded.
44
- - Prefer adding cohesive helper functions here only when behavior is reused across modules.
45
-
46
- ## Verification
47
-
48
- - Run targeted tests for changed helper behavior; run security regressions for auth, filesystem, WebSocket, tunnel, upload, or secret-handling helpers.
49
- - No direct test reference was found by name search; choose the nearest behavioral test or perform a focused smoke check.
50
-
51
- ## Child DOX Index
52
-
53
- No child DOX files.
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:8872cd81aa6c11ea3961f8dc1f085f7518576404dbfb3ad8bc7de3d8dbb58130
3
+ size 3680
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/messages.py CHANGED
@@ -1,75 +1,3 @@
1
- # from . import files
2
-
3
- import json
4
-
5
-
6
- def truncate_text(agent, output, threshold=1000):
7
- threshold = int(threshold)
8
- if not threshold or len(output) <= threshold:
9
- return output
10
-
11
- # Adjust the file path as needed
12
- placeholder = agent.read_prompt(
13
- "fw.msg_truncated.md", length=(len(output) - threshold)
14
- )
15
- # placeholder = files.read_file("./prompts/default/fw.msg_truncated.md", length=(len(output) - threshold))
16
-
17
- start_len = (threshold - len(placeholder)) // 2
18
- end_len = threshold - len(placeholder) - start_len
19
-
20
- truncated_output = output[:start_len] + placeholder + output[-end_len:]
21
- return truncated_output
22
-
23
-
24
- def truncate_dict_by_ratio(agent, data: dict|list|str, threshold_chars: int, truncate_to: int):
25
- threshold_chars = int(threshold_chars)
26
- truncate_to = int(truncate_to)
27
-
28
- def process_item(item):
29
- if isinstance(item, dict):
30
- truncated_dict = {}
31
- cumulative_size = 0
32
-
33
- for key, value in item.items():
34
- processed_value = process_item(value)
35
- serialized_value = json.dumps(processed_value, ensure_ascii=False)
36
- size = len(serialized_value)
37
-
38
- if cumulative_size + size > threshold_chars:
39
- truncated_dict[key] = truncate_text(
40
- agent, serialized_value, truncate_to
41
- )
42
- else:
43
- cumulative_size += size
44
- truncated_dict[key] = processed_value
45
-
46
- return truncated_dict
47
-
48
- elif isinstance(item, list):
49
- truncated_list = []
50
- cumulative_size = 0
51
-
52
- for value in item:
53
- processed_value = process_item(value)
54
- serialized_value = json.dumps(processed_value, ensure_ascii=False)
55
- size = len(serialized_value)
56
-
57
- if cumulative_size + size > threshold_chars:
58
- truncated_list.append(
59
- truncate_text(agent, serialized_value, truncate_to)
60
- )
61
- else:
62
- cumulative_size += size
63
- truncated_list.append(processed_value)
64
-
65
- return truncated_list
66
-
67
- elif isinstance(item, str):
68
- if len(item) > threshold_chars:
69
- return truncate_text(agent, item, truncate_to)
70
- return item
71
-
72
- else:
73
- return item
74
-
75
- return process_item(data)
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:85b47588d6dc155be0101dd047a27d00292e513ad2c83369dfda1ea514ee439b
3
+ size 2472
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/messages.py.dox.md CHANGED
@@ -1,50 +1,3 @@
1
- # messages.py DOX
2
-
3
- ## Purpose
4
-
5
- - Own the `messages.py` helper module.
6
- - This module truncates message payloads for display and storage safety.
7
- - Keep this file-level DOX profile synchronized with `messages.py` because this directory is intentionally flat.
8
-
9
- ## Ownership
10
-
11
- - `messages.py` owns the runtime implementation.
12
- - `messages.py.dox.md` owns durable notes about responsibilities, contracts, side effects, and verification for that implementation.
13
- - Top-level functions:
14
- - `truncate_text(agent, output, threshold=...)`
15
- - `truncate_dict_by_ratio(agent, data: dict | list | str, threshold_chars: int, truncate_to: int)`
16
-
17
- ## Runtime Contracts
18
-
19
- - Helper modules own reusable framework APIs and must preserve public callers unless all callers, tests, and docs are updated together.
20
- - Update this file whenever public functions, classes, persistence behavior, path/security assumptions, side effects, or cross-module contracts change.
21
- - Observed side-effect areas: filesystem reads, filesystem writes, settings/state persistence.
22
- - Imported dependency areas include: `json`.
23
-
24
- ## Key Concepts
25
-
26
- - Important called helpers/classes observed in the source: `agent.read_prompt`, `process_item`, `json.dumps`, `truncate_text`.
27
- - Keep request/response, tool, or helper semantics documented here at the same time as source changes.
28
-
29
- ## Work Guidance
30
-
31
- - Preserve public helper APIs used by core code and plugins unless every caller is updated.
32
- - Keep path, auth, secret, persistence, network, and subprocess behavior explicit and bounded.
33
- - Prefer adding cohesive helper functions here only when behavior is reused across modules.
34
-
35
- ## Verification
36
-
37
- - Run targeted tests for changed helper behavior; run security regressions for auth, filesystem, WebSocket, tunnel, upload, or secret-handling helpers.
38
- - Related tests observed by source search:
39
- - `tests/email_parser_test.py`
40
- - `tests/test_browser_agent_regressions.py`
41
- - `tests/test_chat_compaction.py`
42
- - `tests/test_document_query_fallback.py`
43
- - `tests/test_download_toast_regressions.py`
44
- - `tests/test_fasta2a_client.py`
45
- - `tests/test_mcp_handler_multimodal.py`
46
- - `tests/test_oauth_codex.py`
47
-
48
- ## Child DOX Index
49
-
50
- No child DOX files.
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:2a37e1ed033fccfe1fe484143cb246cafb0278cbc01dafb88a5f7246ddce2db4
3
+ size 2190
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/microsoft_tunnel.py CHANGED
@@ -1,144 +1,3 @@
1
- import getpass
2
- import hashlib
3
- import os
4
- import socket
5
-
6
- from flaredantic import MicrosoftConfig, MicrosoftTunnel, NotifyEvent
7
-
8
- try:
9
- from flaredantic.core.exceptions import MicrosoftTunnelError
10
- except Exception: # pragma: no cover - keeps tests independent from package internals
11
- MicrosoftTunnelError = RuntimeError
12
-
13
- from helpers import files
14
- from helpers.tunnel_common import FlaredanticTunnelHelper
15
-
16
-
17
- MICROSOFT_TUNNEL_ID_ENV_KEYS = (
18
- "A0_MICROSOFT_DEV_TUNNEL_ID",
19
- "MICROSOFT_DEV_TUNNEL_ID",
20
- )
21
- MICROSOFT_TUNNEL_TIMEOUT = 120
22
-
23
-
24
- def default_microsoft_tunnel_id():
25
- for env_key in MICROSOFT_TUNNEL_ID_ENV_KEYS:
26
- configured = (os.environ.get(env_key) or "").strip()
27
- if configured:
28
- return configured
29
-
30
- seed = "|".join([
31
- getpass.getuser(),
32
- socket.gethostname(),
33
- files.get_abs_path("usr"),
34
- ])
35
- digest = hashlib.sha256(seed.encode("utf-8")).hexdigest()[:10]
36
- return f"agent-zero-{digest}"
37
-
38
-
39
- class AgentZeroMicrosoftTunnel(MicrosoftTunnel):
40
- def notify(self, event, message, data=None):
41
- try:
42
- return super().notify(event, message, data)
43
- except AttributeError:
44
- self.agent_zero_notifications.append({
45
- "event": event.value if hasattr(event, "value") else event,
46
- "message": message,
47
- "data": data,
48
- })
49
- return None
50
-
51
- @property
52
- def agent_zero_notifications(self):
53
- if not hasattr(self, "_agent_zero_notifications"):
54
- self._agent_zero_notifications = []
55
- return self._agent_zero_notifications
56
-
57
- def _notify_progress(self, message, data=None):
58
- self.notify(NotifyEvent.INFO, message, data)
59
-
60
- def _ensure_logged_in(self):
61
- parent = getattr(super(), "_ensure_logged_in", None)
62
- if callable(parent):
63
- parent()
64
- self._notify_progress(
65
- "Microsoft Dev Tunnels login confirmed. Preparing your tunnel..."
66
- )
67
-
68
- def _ensure_tunnel(self):
69
- tunnel_id = self.config.tunnel_id
70
- port = str(self.config.port)
71
-
72
- self._notify_progress(
73
- f"Checking Microsoft Dev Tunnel `{tunnel_id}`...",
74
- {"tunnel_id": tunnel_id},
75
- )
76
- show = self._run_cmd(["show", tunnel_id])
77
- if show.returncode != 0:
78
- self._notify_progress(
79
- f"Creating Microsoft Dev Tunnel `{tunnel_id}`...",
80
- {"tunnel_id": tunnel_id},
81
- )
82
- create = self._run_cmd(["create", tunnel_id])
83
- if create.returncode != 0:
84
- raise MicrosoftTunnelError(f"Failed to create tunnel: {create.stdout}")
85
- else:
86
- self._notify_progress(
87
- f"Microsoft Dev Tunnel `{tunnel_id}` already exists. Checking port {port}...",
88
- {"tunnel_id": tunnel_id, "port": port},
89
- )
90
-
91
- self._notify_progress(
92
- f"Checking Microsoft Dev Tunnel port {port}...",
93
- {"tunnel_id": tunnel_id, "port": port},
94
- )
95
- port_show = self._run_cmd(["port", "show", tunnel_id, "-p", port])
96
- if port_show.returncode != 0:
97
- self._notify_progress(
98
- f"Creating Microsoft Dev Tunnel port {port}...",
99
- {"tunnel_id": tunnel_id, "port": port},
100
- )
101
- port_create = self._run_cmd([
102
- "port",
103
- "create",
104
- tunnel_id,
105
- "-p",
106
- port,
107
- "--protocol",
108
- "http",
109
- ])
110
- if port_create.returncode != 0:
111
- raise MicrosoftTunnelError(
112
- f"Failed to create port: {port_create.stdout}"
113
- )
114
-
115
- self._notify_progress(
116
- "Microsoft Dev Tunnel setup is ready. Starting the secure host..."
117
- )
118
-
119
-
120
- class MicrosoftDevTunnel(FlaredanticTunnelHelper):
121
- label = "Microsoft Dev Tunnels"
122
-
123
- def build_tunnel(self):
124
- config = MicrosoftConfig(
125
- port=self.port,
126
- verbose=True,
127
- timeout=MICROSOFT_TUNNEL_TIMEOUT,
128
- tunnel_id=default_microsoft_tunnel_id(),
129
- )
130
- return AgentZeroMicrosoftTunnel(config)
131
-
132
- def start(self):
133
- try:
134
- return super().start()
135
- except Exception as e:
136
- if "Timeout waiting for Microsoft Dev Tunnels URL" not in str(e):
137
- raise
138
- tunnel_id = default_microsoft_tunnel_id()
139
- raise RuntimeError(
140
- "Microsoft Dev Tunnels did not return a URL. Agent Zero uses "
141
- f"the tunnel id `{tunnel_id}` to avoid flaredantic's global "
142
- "`flaredantic` tunnel-id collision. If this still fails, set "
143
- "`A0_MICROSOFT_DEV_TUNNEL_ID` to a fresh unique value and try again."
144
- ) from e
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:7da436b3e30d6d203584c0b38311d2fb262dff93cf346eee7b53103984e6fa0d
3
+ size 4824
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/microsoft_tunnel.py.dox.md CHANGED
@@ -1,50 +1,3 @@
1
- # microsoft_tunnel.py DOX
2
-
3
- ## Purpose
4
-
5
- - Own the `microsoft_tunnel.py` helper module.
6
- - This module implements Microsoft Dev Tunnel provider behavior.
7
- - Keep this file-level DOX profile synchronized with `microsoft_tunnel.py` because this directory is intentionally flat.
8
-
9
- ## Ownership
10
-
11
- - `microsoft_tunnel.py` owns the runtime implementation.
12
- - `microsoft_tunnel.py.dox.md` owns durable notes about responsibilities, contracts, side effects, and verification for that implementation.
13
- - Classes:
14
- - `AgentZeroMicrosoftTunnel` (`MicrosoftTunnel`)
15
- - `notify(self, event, message, data=...)`
16
- - `agent_zero_notifications(self)`
17
- - `MicrosoftDevTunnel` (`FlaredanticTunnelHelper`)
18
- - `build_tunnel(self)`
19
- - `start(self)`
20
- - Top-level functions:
21
- - `default_microsoft_tunnel_id()`
22
- - Notable constants/configuration names: `MICROSOFT_TUNNEL_ID_ENV_KEYS`, `MICROSOFT_TUNNEL_TIMEOUT`.
23
-
24
- ## Runtime Contracts
25
-
26
- - Helper modules own reusable framework APIs and must preserve public callers unless all callers, tests, and docs are updated together.
27
- - Update this file whenever public functions, classes, persistence behavior, path/security assumptions, side effects, or cross-module contracts change.
28
- - Observed side-effect areas: filesystem reads, network calls, settings/state persistence, tunnel state.
29
- - Imported dependency areas include: `flaredantic`, `getpass`, `hashlib`, `helpers`, `helpers.tunnel_common`, `os`, `socket`.
30
-
31
- ## Key Concepts
32
-
33
- - Important called helpers/classes observed in the source: `join`, `strip`, `hashlib.sha256.hexdigest`, `self.notify`, `callable`, `self._notify_progress`, `self._run_cmd`, `MicrosoftConfig`, `AgentZeroMicrosoftTunnel`, `getpass.getuser`, `socket.gethostname`, `files.get_abs_path`, `super.notify`, `parent`, `super.start`, `hashlib.sha256`, `MicrosoftTunnelError`, `default_microsoft_tunnel_id`, `RuntimeError`, `seed.encode`.
34
- - Keep request/response, tool, or helper semantics documented here at the same time as source changes.
35
-
36
- ## Work Guidance
37
-
38
- - Preserve public helper APIs used by core code and plugins unless every caller is updated.
39
- - Keep path, auth, secret, persistence, network, and subprocess behavior explicit and bounded.
40
- - Prefer adding cohesive helper functions here only when behavior is reused across modules.
41
-
42
- ## Verification
43
-
44
- - Run targeted tests for changed helper behavior; run security regressions for auth, filesystem, WebSocket, tunnel, upload, or secret-handling helpers.
45
- - Related tests observed by source search:
46
- - `tests/test_tunnel_remote_link.py`
47
-
48
- ## Child DOX Index
49
-
50
- No child DOX files.
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:4c9da2c11e435e32d7e35482a386bebc99892fed6bdf124e9b1ae52d25749acd
3
+ size 2561
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/migration.py CHANGED
@@ -1,139 +1,3 @@
1
- import os
2
- import json
3
- from helpers import files
4
- from helpers import subagents, extension
5
- from helpers import yaml as yaml_helper
6
- from helpers.print_style import PrintStyle
7
-
8
-
9
- def startup_migration() -> None:
10
- migrate_user_data()
11
- convert_agents_json_yaml()
12
-
13
- extension.call_extensions_sync("startup_migration", None)
14
-
15
-
16
-
17
- def migrate_user_data() -> None:
18
- """
19
- Migrate user data from /tmp and other locations to /usr.
20
- """
21
-
22
- # --- Migrate Directories -------------------------------------------------------
23
- # Move directories from tmp/ or other source locations to usr/
24
-
25
- _move_dir("tmp/chats", "usr/chats")
26
- _move_dir("tmp/scheduler", "usr/scheduler", overwrite=True)
27
- _move_dir("tmp/uploads", "usr/uploads")
28
- _move_dir("tmp/upload", "usr/upload")
29
- _move_dir("tmp/downloads", "usr/downloads")
30
- _move_dir("tmp/email", "usr/email")
31
- _move_dir("knowledge/custom", "usr/knowledge", overwrite=True)
32
-
33
- # --- Migrate Files -------------------------------------------------------------
34
- # Move specific configuration files to usr/
35
-
36
- _move_file("tmp/settings.json", "usr/settings.json")
37
- _move_file("tmp/secrets.env", "usr/secrets.env")
38
- _move_file(".env", "usr/.env", overwrite=True)
39
-
40
- # --- Special Migration Cases ---------------------------------------------------
41
-
42
- # Migrate Memory
43
- _migrate_memory()
44
-
45
- # Flatten default directories (knowledge/default -> knowledge/, etc.)
46
- # We use _merge_dir_contents because we want to move the *contents* of default/
47
- # into the parent directory, not move the default directory itself.
48
- _merge_dir_contents("knowledge/default", "knowledge")
49
-
50
- # --- Cleanup -------------------------------------------------------------------
51
-
52
- # Remove obsolete directories after migration
53
- _cleanup_obsolete()
54
-
55
-
56
- def convert_agents_json_yaml() -> None:
57
- for root in subagents.get_agents_roots():
58
- rel_root = files.deabsolute_path(root)
59
- for subdir in files.get_subdirectories(rel_root):
60
- agent_yaml = os.path.join(rel_root, subdir, "agent.yaml")
61
- if files.exists(agent_yaml):
62
- continue
63
-
64
- agent_json = os.path.join(rel_root, subdir, "agent.json")
65
- if not files.exists(agent_json):
66
- continue
67
-
68
- try:
69
- agent_obj = json.loads(files.read_file(agent_json))
70
- files.write_file(agent_yaml, yaml_helper.dumps(agent_obj))
71
- except Exception as e:
72
- PrintStyle.error(f"Failed to convert {agent_json} to YAML", e)
73
- continue
74
-
75
- # --- Helper Functions ----------------------------------------------------------
76
-
77
- def _move_dir(src: str, dst: str, overwrite: bool = False) -> None:
78
- """
79
- Move a directory from src to dst if src exists and dst does not.
80
- """
81
- if files.exists(src) and (not files.exists(dst) or overwrite):
82
- PrintStyle().print(f"Migrating {src} to {dst}...")
83
- if overwrite and files.exists(dst):
84
- files.delete_dir(dst)
85
- files.move_dir(src, dst)
86
-
87
- def _move_file(src: str, dst: str, overwrite: bool = False) -> None:
88
- """
89
- Move a file from src to dst if src exists and dst does not.
90
- """
91
- if files.exists(src) and (not files.exists(dst) or overwrite):
92
- PrintStyle().print(f"Migrating {src} to {dst}...")
93
- files.move_file(src, dst)
94
-
95
- def _migrate_memory(base_path: str = "memory") -> None:
96
- """
97
- Migrate memory subdirectories.
98
- """
99
- subdirs = files.get_subdirectories(base_path)
100
- for subdir in subdirs:
101
- if subdir == "embeddings":
102
- # Special case: Embeddings
103
- _move_dir("memory/embeddings", "tmp/memory/embeddings")
104
- else:
105
- # Move other memory items to usr/memory
106
- dst = f"usr/memory/{subdir}"
107
- _move_dir(f"memory/{subdir}", dst)
108
-
109
- def _merge_dir_contents(src_parent: str, dst_parent: str) -> None:
110
- """
111
- Moves all items from src_parent to dst_parent.
112
- Useful for flattening structures like 'knowledge/default/*' -> 'knowledge/*'.
113
- """
114
- if not files.exists(src_parent):
115
- return
116
-
117
- entries = files.list_files(src_parent)
118
- for entry in entries:
119
- src = f"{src_parent}/{entry}"
120
- dst = f"{dst_parent}/{entry}"
121
- abs_src = files.get_abs_path(src)
122
- if os.path.isdir(abs_src):
123
- _move_dir(src, dst)
124
- elif os.path.isfile(abs_src):
125
- _move_file(src, dst)
126
-
127
- def _cleanup_obsolete() -> None:
128
- """
129
- Remove directories that are no longer needed.
130
- """
131
- to_remove = [
132
- "knowledge/default",
133
- "memory",
134
- "logs",
135
- ]
136
- for path in to_remove:
137
- if files.exists(path):
138
- PrintStyle().print(f"Removing {path}...")
139
- files.delete_dir(path)
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:32fe773b85f182cf48e2e27d272782460dfe5acb29bc67cad6da0d2c85905e89
3
+ size 4798
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/migration.py.dox.md CHANGED
@@ -1,54 +1,3 @@
1
- # migration.py DOX
2
-
3
- ## Purpose
4
-
5
- - Own the `migration.py` helper module.
6
- - This module migrates legacy user data and runtime layout at startup.
7
- - Keep this file-level DOX profile synchronized with `migration.py` because this directory is intentionally flat.
8
-
9
- ## Ownership
10
-
11
- - `migration.py` owns the runtime implementation.
12
- - `migration.py.dox.md` owns durable notes about responsibilities, contracts, side effects, and verification for that implementation.
13
- - Top-level functions:
14
- - `startup_migration() -> None`
15
- - `migrate_user_data() -> None`: Migrate user data from /tmp and other locations to /usr.
16
- - `convert_agents_json_yaml() -> None`
17
- - `_move_dir(src: str, dst: str, overwrite: bool=...) -> None`: Move a directory from src to dst if src exists and dst does not.
18
- - `_move_file(src: str, dst: str, overwrite: bool=...) -> None`: Move a file from src to dst if src exists and dst does not.
19
- - `_migrate_memory(base_path: str=...) -> None`: Migrate memory subdirectories.
20
- - `_merge_dir_contents(src_parent: str, dst_parent: str) -> None`: Moves all items from src_parent to dst_parent.
21
- - `_cleanup_obsolete() -> None`: Remove directories that are no longer needed.
22
-
23
- ## Runtime Contracts
24
-
25
- - Helper modules own reusable framework APIs and must preserve public callers unless all callers, tests, and docs are updated together.
26
- - Update this file whenever public functions, classes, persistence behavior, path/security assumptions, side effects, or cross-module contracts change.
27
- - Startup cleanup removes obsolete `memory`, `knowledge/default`, and legacy `logs` directories; repeated runs are safe.
28
- - Observed side-effect areas: filesystem reads, filesystem writes, filesystem deletion, settings/state persistence, secret handling, scheduler state.
29
- - Imported dependency areas include: `helpers`, `helpers.print_style`, `json`, `os`.
30
-
31
- ## Key Concepts
32
-
33
- - Important called helpers/classes observed in the source: `migrate_user_data`, `convert_agents_json_yaml`, `extension.call_extensions_sync`, `_move_dir`, `_move_file`, `_migrate_memory`, `_merge_dir_contents`, `_cleanup_obsolete`, `subagents.get_agents_roots`, `files.get_subdirectories`, `files.list_files`, `files.deabsolute_path`, `files.exists`, `files.move_dir`, `files.move_file`, `files.get_abs_path`, `os.path.isdir`, `os.path.join`, `files.delete_dir`, `os.path.isfile`.
34
- - Keep request/response, tool, or helper semantics documented here at the same time as source changes.
35
-
36
- ## Work Guidance
37
-
38
- - Preserve public helper APIs used by core code and plugins unless every caller is updated.
39
- - Keep path, auth, secret, persistence, network, and subprocess behavior explicit and bounded.
40
- - Prefer adding cohesive helper functions here only when behavior is reused across modules.
41
-
42
- ## Verification
43
-
44
- - Run targeted tests for changed helper behavior; run security regressions for auth, filesystem, WebSocket, tunnel, upload, or secret-handling helpers.
45
- - Related tests observed by source search:
46
- - `tests/test_migration_cleanup.py`
47
- - `tests/test_browser_agent_regressions.py`
48
- - `tests/test_office_canvas_setup.py`
49
- - `tests/test_office_document_store.py`
50
- - `tests/test_speech_plugin_split.py`
51
-
52
- ## Child DOX Index
53
-
54
- No child DOX files.
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:5a111e1b05cb3490e7e6156ea4957ea7cf4a48565464eb61329724fec1e1094a
3
+ size 3194
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/modules.py CHANGED
@@ -1,97 +1,3 @@
1
-
2
- import re, os, importlib, importlib.util, inspect, sys
3
- from types import ModuleType
4
- from typing import Any, Type, TypeVar
5
- from helpers.files import get_abs_path
6
- from fnmatch import fnmatch
7
-
8
-
9
- T = TypeVar("T") # Define a generic type variable
10
-
11
-
12
- def import_module(file_path: str) -> ModuleType:
13
- # Handle file paths with periods in the name using importlib.util
14
- abs_path = get_abs_path(file_path)
15
- module_name = os.path.basename(abs_path).replace(".py", "")
16
-
17
- # Create the module spec and load the module
18
- spec = importlib.util.spec_from_file_location(module_name, abs_path)
19
- if spec is None or spec.loader is None:
20
- raise ImportError(f"Could not load module from {abs_path}")
21
-
22
- module = importlib.util.module_from_spec(spec)
23
- spec.loader.exec_module(module)
24
- return module
25
-
26
-
27
- def load_classes_from_folder(
28
- folder: str, name_pattern: str, base_class: Type[T], one_per_file: bool = True
29
- ) -> list[Type[T]]:
30
- classes = []
31
- abs_folder = get_abs_path(folder)
32
-
33
- # Get all .py files in the folder that match the pattern, sorted alphabetically
34
- py_files = sorted(
35
- [
36
- file_name
37
- for file_name in os.listdir(abs_folder)
38
- if fnmatch(file_name, name_pattern) and file_name.endswith(".py")
39
- ]
40
- )
41
-
42
- # Iterate through the sorted list of files
43
- for file_name in py_files:
44
- file_path = os.path.join(abs_folder, file_name)
45
- # Use the new import_module function
46
- module = import_module(file_path)
47
-
48
- # Get all classes in the module
49
- class_list = inspect.getmembers(module, inspect.isclass)
50
-
51
- # Filter for classes that are subclasses of the given base_class
52
- # iterate backwards to skip imported superclasses
53
- for cls in reversed(class_list):
54
- if cls[1] is not base_class and issubclass(cls[1], base_class):
55
- classes.append(cls[1])
56
- if one_per_file:
57
- break
58
-
59
- return classes
60
-
61
-
62
- def load_classes_from_file(
63
- file: str, base_class: type[T], one_per_file: bool = True
64
- ) -> list[type[T]]:
65
- classes = []
66
- # Use the new import_module function
67
- module = import_module(file)
68
-
69
- # Get all classes in the module
70
- class_list = inspect.getmembers(module, inspect.isclass)
71
-
72
- # Filter for classes that are subclasses of the given base_class
73
- # iterate backwards to skip imported superclasses
74
- for cls in reversed(class_list):
75
- if cls[1] is not base_class and issubclass(cls[1], base_class):
76
- classes.append(cls[1])
77
- if one_per_file:
78
- break
79
-
80
- return classes
81
-
82
-
83
- def purge_namespace(namespace: str):
84
- to_delete = [
85
- name
86
- for name in sys.modules
87
- if name == namespace or name.startswith(namespace + ".")
88
- ]
89
-
90
- # delete deepest first just to be tidy
91
- to_delete.sort(key=lambda n: n.count("."), reverse=True)
92
-
93
- for name in to_delete:
94
- del sys.modules[name]
95
-
96
- importlib.invalidate_caches()
97
- return to_delete
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:0c70c8a4e53bde7172a8e5a3e38ae91c0bd3b086c013f0c7c622724f0b59e4cc
3
+ size 3010
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/modules.py.dox.md CHANGED
@@ -1,53 +1,3 @@
1
- # modules.py DOX
2
-
3
- ## Purpose
4
-
5
- - Own the `modules.py` helper module.
6
- - This module imports modules and classes dynamically with namespace cache control.
7
- - Keep this file-level DOX profile synchronized with `modules.py` because this directory is intentionally flat.
8
-
9
- ## Ownership
10
-
11
- - `modules.py` owns the runtime implementation.
12
- - `modules.py.dox.md` owns durable notes about responsibilities, contracts, side effects, and verification for that implementation.
13
- - Top-level functions:
14
- - `import_module(file_path: str) -> ModuleType`
15
- - `load_classes_from_folder(folder: str, name_pattern: str, base_class: Type[T], one_per_file: bool=...) -> list[Type[T]]`
16
- - `load_classes_from_file(file: str, base_class: type[T], one_per_file: bool=...) -> list[type[T]]`
17
- - `purge_namespace(namespace: str)`
18
- - Notable constants/configuration names: `T`.
19
-
20
- ## Runtime Contracts
21
-
22
- - Helper modules own reusable framework APIs and must preserve public callers unless all callers, tests, and docs are updated together.
23
- - Update this file whenever public functions, classes, persistence behavior, path/security assumptions, side effects, or cross-module contracts change.
24
- - Observed side-effect areas: filesystem reads, filesystem deletion.
25
- - Imported dependency areas include: `fnmatch`, `helpers.files`, `importlib`, `importlib.util`, `inspect`, `os`, `re`, `sys`, `types`, `typing`.
26
-
27
- ## Key Concepts
28
-
29
- - Important called helpers/classes observed in the source: `TypeVar`, `get_abs_path`, `os.path.basename.replace`, `importlib.util.spec_from_file_location`, `importlib.util.module_from_spec`, `spec.loader.exec_module`, `import_module`, `inspect.getmembers`, `to_delete.sort`, `importlib.invalidate_caches`, `ImportError`, `os.path.join`, `os.path.basename`, `issubclass`, `os.listdir`, `name.startswith`, `n.count`, `fnmatch`, `file_name.endswith`.
30
- - Keep request/response, tool, or helper semantics documented here at the same time as source changes.
31
-
32
- ## Work Guidance
33
-
34
- - Preserve public helper APIs used by core code and plugins unless every caller is updated.
35
- - Keep path, auth, secret, persistence, network, and subprocess behavior explicit and bounded.
36
- - Prefer adding cohesive helper functions here only when behavior is reused across modules.
37
-
38
- ## Verification
39
-
40
- - Run targeted tests for changed helper behavior; run security regressions for auth, filesystem, WebSocket, tunnel, upload, or secret-handling helpers.
41
- - Related tests observed by source search:
42
- - `tests/test_a0_connector_prompt_gating.py`
43
- - `tests/test_browser_agent_regressions.py`
44
- - `tests/test_docker_release_plan.py`
45
- - `tests/test_error_retry_plugin.py`
46
- - `tests/test_file_tree_visualize.py`
47
- - `tests/test_git_version_label.py`
48
- - `tests/test_mcp_handler_multimodal.py`
49
- - `tests/test_memory_quality.py`
50
-
51
- ## Child DOX Index
52
-
53
- No child DOX files.
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:10b2bb916631cb41b19aee53a407783a4d5650b2d33e07d73b780d4be61ebb10
3
+ size 2809
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/network.py CHANGED
@@ -1,202 +1,3 @@
1
- from __future__ import annotations
2
-
3
- from dataclasses import dataclass
4
- import ipaddress
5
- import os
6
- import socket
7
- import struct
8
- from urllib.parse import urljoin, urlparse
9
-
10
- import requests
11
-
12
-
13
- SAFE_HTTP_SCHEMES = frozenset({"http", "https"})
14
- DEFAULT_FETCH_TIMEOUT = (3.05, 10.0)
15
- DEFAULT_HTTP_USER_AGENT = "@mixedbread-ai/unstructured"
16
-
17
-
18
- @dataclass(frozen=True)
19
- class HttpFetchResult:
20
- url: str
21
- content: bytes
22
- content_type: str | None
23
- encoding: str | None
24
-
25
-
26
- class UnsafeUrlError(ValueError):
27
- """Raised when a remote URL resolves to a non-public destination."""
28
-
29
-
30
- def _build_request_headers() -> dict[str, str]:
31
- user_agent = (
32
- os.getenv("USER_AGENT")
33
- or os.getenv("user_agent")
34
- or DEFAULT_HTTP_USER_AGENT
35
- ).strip()
36
- return {"User-Agent": user_agent or DEFAULT_HTTP_USER_AGENT}
37
-
38
-
39
- def _normalize_content_type(content_type: str | None) -> str | None:
40
- if not content_type:
41
- return None
42
- return content_type.split(";", 1)[0].strip().lower() or None
43
-
44
-
45
- def resolve_host_ips(hostname: str) -> tuple[ipaddress._BaseAddress, ...]:
46
- try:
47
- results = socket.getaddrinfo(
48
- hostname,
49
- None,
50
- family=socket.AF_UNSPEC,
51
- type=socket.SOCK_STREAM,
52
- )
53
- except socket.gaierror as exc:
54
- raise UnsafeUrlError(f"Unable to resolve hostname '{hostname}'") from exc
55
-
56
- ips: list[ipaddress._BaseAddress] = []
57
- seen: set[str] = set()
58
- for _family, _type, _proto, _canonname, sockaddr in results:
59
- address = sockaddr[0]
60
- if "%" in address:
61
- address = address.split("%", 1)[0]
62
- ip = ipaddress.ip_address(address)
63
- key = ip.compressed
64
- if key in seen:
65
- continue
66
- seen.add(key)
67
- ips.append(ip)
68
-
69
- if not ips:
70
- raise UnsafeUrlError(f"Hostname '{hostname}' did not resolve to an IP address")
71
-
72
- return tuple(ips)
73
-
74
-
75
- def validate_public_http_url(url: str) -> tuple[ipaddress._BaseAddress, ...]:
76
- parsed = urlparse(url)
77
-
78
- if parsed.scheme not in SAFE_HTTP_SCHEMES:
79
- raise UnsafeUrlError("Only http:// and https:// URLs are supported")
80
- if not parsed.hostname:
81
- raise UnsafeUrlError("URL hostname is required")
82
- if parsed.username or parsed.password:
83
- raise UnsafeUrlError("URLs with embedded credentials are not allowed")
84
-
85
- hostname = parsed.hostname.rstrip(".").lower()
86
- if hostname == "localhost" or hostname.endswith(".localhost"):
87
- raise UnsafeUrlError(f"Blocked local hostname '{hostname}'")
88
-
89
- ips = resolve_host_ips(hostname)
90
- blocked = [str(ip) for ip in ips if not ip.is_global]
91
- if blocked:
92
- raise UnsafeUrlError(
93
- f"Blocked non-public address resolution for '{hostname}': {', '.join(blocked)}"
94
- )
95
-
96
- return ips
97
-
98
-
99
- def fetch_public_http_resource(
100
- url: str,
101
- *,
102
- max_bytes: int,
103
- max_redirects: int = 5,
104
- timeout: tuple[float, float] = DEFAULT_FETCH_TIMEOUT,
105
- ) -> HttpFetchResult:
106
- current_url = url
107
- session = requests.Session()
108
- session.trust_env = False
109
-
110
- for redirect_count in range(max_redirects + 1):
111
- validate_public_http_url(current_url)
112
-
113
- try:
114
- with session.get(
115
- current_url,
116
- stream=True,
117
- allow_redirects=False,
118
- headers=_build_request_headers(),
119
- timeout=timeout,
120
- ) as response:
121
- if 300 <= response.status_code < 400:
122
- location = response.headers.get("Location")
123
- if not location:
124
- raise ValueError(
125
- f"Remote URL redirect is missing a Location header: {current_url}"
126
- )
127
- if redirect_count >= max_redirects:
128
- raise ValueError(
129
- f"Remote URL exceeded redirect limit ({max_redirects}): {url}"
130
- )
131
- current_url = urljoin(current_url, location)
132
- continue
133
-
134
- if response.status_code >= 400:
135
- raise ValueError(
136
- f"Remote URL returned HTTP {response.status_code}: {current_url}"
137
- )
138
-
139
- content_length = response.headers.get("Content-Length")
140
- if content_length:
141
- try:
142
- declared_length = int(content_length)
143
- except ValueError:
144
- declared_length = None
145
- if declared_length is not None and declared_length > max_bytes:
146
- raise ValueError(
147
- f"Remote document exceeds max size {max_bytes} bytes: {current_url}"
148
- )
149
-
150
- body = bytearray()
151
- for chunk in response.iter_content(chunk_size=64 * 1024):
152
- if not chunk:
153
- continue
154
- body.extend(chunk)
155
- if len(body) > max_bytes:
156
- raise ValueError(
157
- f"Remote document exceeds max size {max_bytes} bytes: {current_url}"
158
- )
159
-
160
- return HttpFetchResult(
161
- url=current_url,
162
- content=bytes(body),
163
- content_type=_normalize_content_type(
164
- response.headers.get("Content-Type")
165
- ),
166
- encoding=response.encoding,
167
- )
168
- except requests.RequestException as exc:
169
- raise ValueError(
170
- f"Remote document fetch failed for {current_url}: {exc}"
171
- ) from exc
172
-
173
- raise ValueError(f"Remote URL exceeded redirect limit ({max_redirects}): {url}")
174
-
175
-
176
- def is_loopback_address(address: str) -> bool:
177
- """Check whether *address* resolves to a loopback interface."""
178
- _checkers = {
179
- socket.AF_INET: lambda x: (
180
- struct.unpack("!I", socket.inet_aton(x))[0] >> (32 - 8)
181
- ) == 127,
182
- socket.AF_INET6: lambda x: x == "::1",
183
- }
184
- try:
185
- socket.inet_pton(socket.AF_INET6, address)
186
- return _checkers[socket.AF_INET6](address)
187
- except socket.error:
188
- pass
189
- try:
190
- socket.inet_pton(socket.AF_INET, address)
191
- return _checkers[socket.AF_INET](address)
192
- except socket.error:
193
- pass
194
- for family in (socket.AF_INET, socket.AF_INET6):
195
- try:
196
- r = socket.getaddrinfo(address, None, family, socket.SOCK_STREAM)
197
- except socket.gaierror:
198
- return False
199
- for fam, _, _, _, sockaddr in r:
200
- if not _checkers[fam](sockaddr[0]):
201
- return False
202
- return True
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:a6083e43f235ecd1470957f3fc33999150bd018bf2b4944a46dcb8a74474b0a1
3
+ size 6715
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/network.py.dox.md CHANGED
@@ -1,55 +1,3 @@
1
- # network.py DOX
2
-
3
- ## Purpose
4
-
5
- - Own the `network.py` helper module.
6
- - This module validates public URLs and fetches remote resources safely.
7
- - Keep this file-level DOX profile synchronized with `network.py` because this directory is intentionally flat.
8
-
9
- ## Ownership
10
-
11
- - `network.py` owns the runtime implementation.
12
- - `network.py.dox.md` owns durable notes about responsibilities, contracts, side effects, and verification for that implementation.
13
- - Classes:
14
- - `HttpFetchResult` (no explicit base class)
15
- - `UnsafeUrlError` (`ValueError`)
16
- - Top-level functions:
17
- - `_build_request_headers() -> dict[str, str]`
18
- - `_normalize_content_type(content_type: str | None) -> str | None`
19
- - `resolve_host_ips(hostname: str) -> tuple[ipaddress._BaseAddress, ...]`
20
- - `validate_public_http_url(url: str) -> tuple[ipaddress._BaseAddress, ...]`
21
- - `fetch_public_http_resource(url: str, max_bytes: int, max_redirects: int=..., timeout: tuple[float, float]=...) -> HttpFetchResult`
22
- - `is_loopback_address(address: str) -> bool`: Check whether *address* resolves to a loopback interface.
23
- - Notable constants/configuration names: `SAFE_HTTP_SCHEMES`, `DEFAULT_FETCH_TIMEOUT`, `DEFAULT_HTTP_USER_AGENT`.
24
-
25
- ## Runtime Contracts
26
-
27
- - Helper modules own reusable framework APIs and must preserve public callers unless all callers, tests, and docs are updated together.
28
- - Update this file whenever public functions, classes, persistence behavior, path/security assumptions, side effects, or cross-module contracts change.
29
- - Observed side-effect areas: network calls, secret handling.
30
- - Imported dependency areas include: `__future__`, `dataclasses`, `ipaddress`, `os`, `requests`, `socket`, `struct`, `urllib.parse`.
31
-
32
- ## Key Concepts
33
-
34
- - Important called helpers/classes observed in the source: `frozenset`, `dataclass`, `strip`, `urlparse`, `parsed.hostname.rstrip.lower`, `resolve_host_ips`, `requests.Session`, `ValueError`, `content_type.split.strip.lower`, `socket.getaddrinfo`, `ipaddress.ip_address`, `seen.add`, `UnsafeUrlError`, `hostname.endswith`, `validate_public_http_url`, `socket.inet_pton`, `_checkers`, `parsed.hostname.rstrip`, `os.getenv`, `content_type.split.strip`.
35
- - Keep request/response, tool, or helper semantics documented here at the same time as source changes.
36
-
37
- ## Work Guidance
38
-
39
- - Preserve public helper APIs used by core code and plugins unless every caller is updated.
40
- - Keep path, auth, secret, persistence, network, and subprocess behavior explicit and bounded.
41
- - Prefer adding cohesive helper functions here only when behavior is reused across modules.
42
-
43
- ## Verification
44
-
45
- - Run targeted tests for changed helper behavior; run security regressions for auth, filesystem, WebSocket, tunnel, upload, or secret-handling helpers.
46
- - Related tests observed by source search:
47
- - `tests/test_oauth_gemini_api.py`
48
- - `tests/test_oauth_github_copilot.py`
49
- - `tests/test_oauth_xai_grok.py`
50
- - `tests/test_plugin_scan_prompt.py`
51
- - `tests/test_tunnel_remote_link.py`
52
-
53
- ## Child DOX Index
54
-
55
- No child DOX files.
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:10f8c25855e46830f82f4ddc219f032613274bd2d348fbd4a4faab36d4ca70c7
3
+ size 3001
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/notification.py CHANGED
@@ -1,250 +1,3 @@
1
- from dataclasses import dataclass
2
- import uuid
3
- import threading
4
- from datetime import datetime, timedelta
5
- from enum import Enum
6
-
7
- from helpers.localization import Localization
8
-
9
-
10
- class NotificationType(Enum):
11
- INFO = "info"
12
- SUCCESS = "success"
13
- WARNING = "warning"
14
- ERROR = "error"
15
- PROGRESS = "progress"
16
-
17
-
18
- class NotificationPriority(Enum):
19
- NORMAL = 10
20
- HIGH = 20
21
-
22
-
23
- @dataclass
24
- class NotificationItem:
25
- manager: "NotificationManager"
26
- no: int
27
- type: NotificationType
28
- priority: NotificationPriority
29
- title: str
30
- message: str
31
- detail: str # HTML content for expandable details
32
- timestamp: datetime
33
- display_time: int = 3 # Display duration in seconds, default 3 seconds
34
- read: bool = False
35
- id: str = ""
36
- group: str = "" # Group identifier for grouping related notifications
37
-
38
- def __post_init__(self):
39
- if not self.id:
40
- self.id = str(uuid.uuid4())
41
- # Ensure type is always NotificationType
42
- if isinstance(self.type, str):
43
- self.type = NotificationType(self.type)
44
-
45
- def mark_read(self):
46
- self.read = True
47
- self.manager.update_item(self.no, read=True)
48
-
49
- def output(self):
50
- return {
51
- "no": self.no,
52
- "id": self.id,
53
- "type": self.type.value if isinstance(self.type, NotificationType) else self.type,
54
- "priority": self.priority.value if isinstance(self.priority, NotificationPriority) else self.priority,
55
- "title": self.title,
56
- "message": self.message,
57
- "detail": self.detail,
58
- "timestamp": Localization.get().serialize_datetime(self.timestamp),
59
- "display_time": self.display_time,
60
- "read": self.read,
61
- "group": self.group,
62
- }
63
-
64
-
65
- class NotificationManager:
66
- def __init__(self, max_notifications: int = 100):
67
- self._lock = threading.RLock()
68
- self.guid: str = str(uuid.uuid4())
69
- self.updates: list[int] = []
70
- self.notifications: list[NotificationItem] = []
71
- self.max_notifications = max_notifications
72
-
73
- @staticmethod
74
- def send_notification(
75
- type: NotificationType,
76
- priority: NotificationPriority,
77
- message: str,
78
- title: str = "",
79
- detail: str = "",
80
- display_time: int = 3,
81
- group: str = "",
82
- id: str = "",
83
- ) -> NotificationItem:
84
- from agent import AgentContext
85
- return AgentContext.get_notification_manager().add_notification(
86
- type, priority, message, title, detail, display_time, group, id
87
- )
88
-
89
- def add_notification(
90
- self,
91
- type: NotificationType,
92
- priority: NotificationPriority,
93
- message: str,
94
- title: str = "",
95
- detail: str = "",
96
- display_time: int = 3,
97
- group: str = "",
98
- id: str = "",
99
- ) -> NotificationItem:
100
- with self._lock:
101
- existing = None
102
- if id:
103
- existing = next((n for n in self.notifications if n.id == id), None)
104
-
105
- if existing:
106
- existing.type = NotificationType(type)
107
- existing.priority = NotificationPriority(priority)
108
- existing.title = title
109
- existing.message = message
110
- existing.detail = detail
111
- existing.timestamp = Localization.get().now()
112
- existing.display_time = display_time
113
- existing.group = group
114
- existing.read = False
115
- self.updates.append(existing.no)
116
- item = existing
117
- else:
118
- # Create notification item
119
- item = NotificationItem(
120
- manager=self,
121
- no=len(self.notifications),
122
- type=NotificationType(type),
123
- priority=NotificationPriority(priority),
124
- title=title,
125
- message=message,
126
- detail=detail,
127
- timestamp=Localization.get().now(),
128
- display_time=display_time,
129
- id=id,
130
- group=group,
131
- )
132
-
133
- self.notifications.append(item)
134
- self.updates.append(item.no)
135
- self._enforce_limit()
136
-
137
- from helpers.state_monitor_integration import mark_dirty_all
138
- mark_dirty_all(reason="notification.NotificationManager.add_notification")
139
- return item
140
-
141
- def _enforce_limit(self):
142
- with self._lock:
143
- if len(self.notifications) > self.max_notifications:
144
- # Remove oldest notifications
145
- to_remove = len(self.notifications) - self.max_notifications
146
- self.notifications = self.notifications[to_remove:]
147
- # Adjust notification numbers
148
- for i, notification in enumerate(self.notifications):
149
- notification.no = i
150
- # Adjust updates list
151
- self.updates = [no - to_remove for no in self.updates if no >= to_remove]
152
-
153
- def get_recent_notifications(self, seconds: int = 30) -> list[NotificationItem]:
154
- cutoff = Localization.get().now() - timedelta(seconds=seconds)
155
- with self._lock:
156
- return [n for n in self.notifications if n.timestamp >= cutoff]
157
-
158
- def output(self, start: int | None = None, end: int | None = None) -> list[dict]:
159
- return self.output_with_state(start, end)[0]
160
-
161
- def output_with_state(
162
- self, start: int | None = None, end: int | None = None
163
- ) -> tuple[list[dict], str, int]:
164
- with self._lock:
165
- if start is None:
166
- start = 0
167
- if end is None:
168
- end = len(self.updates)
169
- updates = self.updates[start:end]
170
- out = []
171
- seen = set()
172
- for update in updates:
173
- if update not in seen and update < len(self.notifications):
174
- out.append(self.notifications[update].output())
175
- seen.add(update)
176
- return out, self.guid, len(self.updates)
177
-
178
- def output_all(self) -> list[dict]:
179
- with self._lock:
180
- notifications = list(self.notifications)
181
- return [n.output() for n in notifications]
182
-
183
- def mark_read_by_ids(self, notification_ids: list[str]) -> int:
184
- ids = {nid for nid in notification_ids if isinstance(nid, str) and nid.strip()}
185
- if not ids:
186
- return 0
187
-
188
- changed_nos: list[int] = []
189
- with self._lock:
190
- for notification in self.notifications:
191
- if notification.id in ids and not notification.read:
192
- notification.read = True
193
- changed_nos.append(notification.no)
194
- if changed_nos:
195
- self.updates.extend(changed_nos)
196
-
197
- if not changed_nos:
198
- return 0
199
-
200
- from helpers.state_monitor_integration import mark_dirty_all
201
- mark_dirty_all(reason="notification.NotificationManager.mark_read_by_ids")
202
- return len(changed_nos)
203
-
204
- def update_item(self, no: int, **kwargs) -> None:
205
- self._update_item(no, **kwargs)
206
-
207
- def _update_item(self, no: int, **kwargs):
208
- changed = False
209
- with self._lock:
210
- if no < len(self.notifications):
211
- item = self.notifications[no]
212
- for key, value in kwargs.items():
213
- if hasattr(item, key):
214
- setattr(item, key, value)
215
- self.updates.append(no)
216
- changed = True
217
-
218
- if not changed:
219
- return
220
-
221
- from helpers.state_monitor_integration import mark_dirty_all
222
- mark_dirty_all(reason="notification.NotificationManager._update_item")
223
-
224
- def mark_all_read(self):
225
- changed_nos: list[int] = []
226
- with self._lock:
227
- for notification in self.notifications:
228
- if not notification.read:
229
- notification.read = True
230
- changed_nos.append(notification.no)
231
- if changed_nos:
232
- self.updates.extend(changed_nos)
233
-
234
- if not changed_nos:
235
- return
236
-
237
- from helpers.state_monitor_integration import mark_dirty_all
238
- mark_dirty_all(reason="notification.NotificationManager.mark_all_read")
239
-
240
- def clear_all(self):
241
- with self._lock:
242
- self.notifications = []
243
- self.updates = []
244
- self.guid = str(uuid.uuid4())
245
- from helpers.state_monitor_integration import mark_dirty_all
246
- mark_dirty_all(reason="notification.NotificationManager.clear_all")
247
-
248
- def get_notifications_by_type(self, type: NotificationType) -> list[NotificationItem]:
249
- with self._lock:
250
- return [n for n in self.notifications if n.type == type]
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:9807efbceeeb847f65df5b7896318718edeadc10fbc996f0ebc3b4557510e64c
3
+ size 8793
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/notification.py.dox.md CHANGED
@@ -1,64 +1,3 @@
1
- # notification.py DOX
2
-
3
- ## Purpose
4
-
5
- - Own the `notification.py` helper module.
6
- - This module owns notification data models and notification manager state.
7
- - Keep this file-level DOX profile synchronized with `notification.py` because this directory is intentionally flat.
8
-
9
- ## Ownership
10
-
11
- - `notification.py` owns the runtime implementation.
12
- - `notification.py.dox.md` owns durable notes about responsibilities, contracts, side effects, and verification for that implementation.
13
- - Classes:
14
- - `NotificationType` (`Enum`)
15
- - `NotificationPriority` (`Enum`)
16
- - `NotificationItem` (no explicit base class)
17
- - `mark_read(self)`
18
- - `output(self)`
19
- - `NotificationManager` (no explicit base class)
20
- - `send_notification(type: NotificationType, priority: NotificationPriority, message: str, title: str=..., detail: str=..., display_time: int=..., group: str=..., id: str=...) -> NotificationItem`
21
- - `add_notification(self, type: NotificationType, priority: NotificationPriority, message: str, title: str=..., detail: str=..., display_time: int=..., group: str=..., id: str=...) -> NotificationItem`
22
- - `get_recent_notifications(self, seconds: int=...) -> list[NotificationItem]`
23
- - `output(self, start: int | None=..., end: int | None=...) -> list[dict]`
24
- - `output_with_state(self, start: int | None=..., end: int | None=...) -> tuple[list[dict], str, int]`
25
- - `output_all(self) -> list[dict]`
26
- - `mark_read_by_ids(self, notification_ids: list[str]) -> int`
27
- - `update_item(self, no: int, **kwargs) -> None`
28
- - `mark_all_read(self)`
29
-
30
- ## Runtime Contracts
31
-
32
- - Helper modules own reusable framework APIs and must preserve public callers unless all callers, tests, and docs are updated together.
33
- - Update this file whenever public functions, classes, persistence behavior, path/security assumptions, side effects, or cross-module contracts change.
34
- - Notification payloads, GUIDs, and update cursors are captured under one lock so WebUI snapshots cannot skip notifications created during snapshot assembly.
35
- - Observed side-effect areas: filesystem deletion, settings/state persistence.
36
- - Imported dependency areas include: `dataclasses`, `datetime`, `enum`, `helpers.localization`, `threading`, `uuid`.
37
-
38
- ## Key Concepts
39
-
40
- - Important called helpers/classes observed in the source: `self.manager.update_item`, `threading.RLock`, `AgentContext.get_notification_manager.add_notification`, `mark_dirty_all`, `self._update_item`, `NotificationType`, `Localization.get.serialize_datetime`, `uuid.uuid4`, `Localization.get.now`, `timedelta`, `n.output`, `AgentContext.get_notification_manager`, `next`, `NotificationPriority`, `NotificationItem`, `self._enforce_limit`, `seen.add`, `notifications.output`, `nid.strip`.
41
- - Keep request/response, tool, or helper semantics documented here at the same time as source changes.
42
-
43
- ## Work Guidance
44
-
45
- - Preserve public helper APIs used by core code and plugins unless every caller is updated.
46
- - Keep path, auth, secret, persistence, network, and subprocess behavior explicit and bounded.
47
- - Prefer adding cohesive helper functions here only when behavior is reused across modules.
48
-
49
- ## Verification
50
-
51
- - Run targeted tests for changed helper behavior; run security regressions for auth, filesystem, WebSocket, tunnel, upload, or secret-handling helpers.
52
- - Related tests observed by source search:
53
- - `tests/test_download_toast_regressions.py`
54
- - `tests/test_multi_tab_isolation.py`
55
- - `tests/test_self_update_tag_filter.py`
56
- - `tests/test_snapshot_parity.py`
57
- - `tests/test_snapshot_schema_v1.py`
58
- - `tests/test_state_monitor.py`
59
- - `tests/test_state_sync_handler.py`
60
- - `tests/test_state_sync_welcome_screen.py`
61
-
62
- ## Child DOX Index
63
-
64
- No child DOX files.
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:65e550506f83a76a4721a8bda8eca3821ef4c5c93294fde195fb5e0fb46cb544
3
+ size 3684
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/parallel_tools.py CHANGED
@@ -1,762 +1,3 @@
1
- from __future__ import annotations
2
-
3
- import asyncio
4
- import json
5
- import time
6
- import uuid
7
- from dataclasses import dataclass, field, replace
8
- from typing import Any, Literal, TYPE_CHECKING
9
-
10
- from helpers import extract_tools
11
- from helpers.defer import DeferredTask, THREAD_BACKGROUND
12
- from helpers.extension import call_extensions_async
13
- from helpers.print_style import PrintStyle
14
-
15
- if TYPE_CHECKING:
16
- from agent import Agent, AgentConfig, AgentContext
17
- from helpers.log import LogItem
18
-
19
-
20
- PARALLEL_JOBS_KEY = "_parallel_jobs"
21
- PARALLEL_WORKER_PARENT_CONTEXT_KEY = "_parallel_parent_context_id"
22
- PARALLEL_WORKER_JOB_KEY = "_parallel_job_id"
23
- PARALLEL_WORKER_KIND_KEY = "_parallel_worker_kind"
24
-
25
- CHILD_PARENT_CONTEXT_ID_KEY = "parent_context_id"
26
- CHILD_PARENT_AGENT_NUMBER_KEY = "parent_agent_number"
27
- CHILD_PARENT_CONTEXT_KIND_KEY = "parent_context_kind"
28
- CHILD_PARENT_CONTEXT_LABEL_KEY = "parent_context_label"
29
- CHILD_PARALLEL_JOB_ID_KEY = "parallel_job_id"
30
- CHILD_PARALLEL_TOOL_NAME_KEY = "parallel_tool_name"
31
-
32
- DEFAULT_MAX_CALLS = 8
33
- DEFAULT_TIMEOUT_SECONDS = 300
34
- POLL_INTERVAL_SECONDS = 0.5
35
- DISALLOWED_PARALLEL_TOOLS = {"document_query", "response"}
36
-
37
- TERMINAL_STATES = {"success", "error", "cancelled", "timeout"}
38
- JobState = Literal["pending", "running", "success", "error", "cancelled", "timeout"]
39
- JobKind = Literal["tool", "subordinate"]
40
-
41
-
42
- @dataclass
43
- class NormalizedToolCall:
44
- index: int
45
- tool_name: str
46
- tool_args: dict[str, Any]
47
-
48
-
49
- @dataclass
50
- class ParallelJob:
51
- id: str
52
- parent_context_id: str
53
- index: int
54
- tool_name: str
55
- tool_args: dict[str, Any]
56
- kind: JobKind
57
- parent_agent: "Agent | None" = field(default=None, repr=False)
58
- state: JobState = "pending"
59
- created_at: float = field(default_factory=time.time)
60
- started_at: float | None = None
61
- completed_at: float | None = None
62
- result: str | None = None
63
- error: str | None = None
64
- worker_context_id: str | None = None
65
- log_id: str = field(default_factory=lambda: str(uuid.uuid4()))
66
- log_item: "LogItem | None" = field(default=None, repr=False)
67
- deferred_task: DeferredTask | None = field(default=None, repr=False)
68
- parent_history: list[tuple[Any, int]] = field(default_factory=list, repr=False)
69
-
70
- def elapsed(self) -> float:
71
- end = self.completed_at or time.time()
72
- start = self.started_at or self.created_at
73
- return max(0.0, end - start)
74
-
75
-
76
- def extract_tool_calls(args: dict[str, Any]) -> Any:
77
- for key in ("tool_calls", "calls", "items"):
78
- if key in args:
79
- return args.get(key)
80
- return None
81
-
82
-
83
- def normalize_parallel_tool_calls(raw_calls: Any) -> list[NormalizedToolCall]:
84
- if isinstance(raw_calls, str):
85
- try:
86
- raw_calls = json.loads(raw_calls)
87
- except json.JSONDecodeError as exc:
88
- raise ValueError(
89
- "`tool_calls` must be an array of normal tool-call objects."
90
- ) from exc
91
- if not isinstance(raw_calls, list):
92
- raise ValueError("`tool_calls` must be an array of normal tool-call objects.")
93
- if not raw_calls:
94
- raise ValueError("`tool_calls` must contain at least one tool call.")
95
- if len(raw_calls) > DEFAULT_MAX_CALLS:
96
- raise ValueError(f"`tool_calls` supports at most {DEFAULT_MAX_CALLS} items.")
97
-
98
- calls: list[NormalizedToolCall] = []
99
- for index, raw_call in enumerate(raw_calls):
100
- try:
101
- tool_name, tool_args = extract_tools.normalize_tool_request(raw_call)
102
- except ValueError as exc:
103
- raise ValueError(f"tool_calls[{index}] is not a valid tool call: {exc}") from exc
104
-
105
- if tool_name == "parallel":
106
- raise ValueError("`parallel` cannot be nested inside another `parallel` call.")
107
- if tool_name in DISALLOWED_PARALLEL_TOOLS:
108
- raise ValueError(
109
- f"`{tool_name}` cannot be used inside `parallel`; call it sequentially."
110
- )
111
-
112
- calls.append(
113
- NormalizedToolCall(
114
- index=index,
115
- tool_name=tool_name,
116
- tool_args=dict(tool_args),
117
- )
118
- )
119
- return calls
120
-
121
-
122
- def normalize_job_ids(raw_job_ids: Any) -> list[str]:
123
- if raw_job_ids is None:
124
- return []
125
- if isinstance(raw_job_ids, str):
126
- return [raw_job_ids]
127
- if isinstance(raw_job_ids, list):
128
- return [str(item) for item in raw_job_ids if str(item).strip()]
129
- raise ValueError("`job_ids` must be a string or an array of strings.")
130
-
131
-
132
- def coerce_bool(value: Any, default: bool) -> bool:
133
- if value is None:
134
- return default
135
- if isinstance(value, bool):
136
- return value
137
- if isinstance(value, str):
138
- normalized = value.strip().lower()
139
- if normalized in {"1", "true", "yes", "on"}:
140
- return True
141
- if normalized in {"0", "false", "no", "off"}:
142
- return False
143
- return bool(value)
144
-
145
-
146
- def coerce_timeout(value: Any) -> int:
147
- if value in (None, ""):
148
- return DEFAULT_TIMEOUT_SECONDS
149
- try:
150
- timeout = int(value)
151
- except (TypeError, ValueError) as exc:
152
- raise ValueError(f"`timeout` must be an integer number of seconds, got {value!r}.") from exc
153
- if timeout <= 0:
154
- raise ValueError("`timeout` must be greater than 0.")
155
- return timeout
156
-
157
-
158
- def _parallel_worker_kind(agent: "Agent | None") -> JobKind | None:
159
- context = getattr(agent, "context", None)
160
- if not context:
161
- return None
162
- kind = context.get_data(PARALLEL_WORKER_KIND_KEY)
163
- if kind in {"tool", "subordinate"}:
164
- return kind
165
- if context.get_data(PARALLEL_WORKER_JOB_KEY):
166
- return "tool"
167
- return None
168
-
169
-
170
- def is_parallel_worker(agent: "Agent | None") -> bool:
171
- return _parallel_worker_kind(agent) == "tool"
172
-
173
-
174
- def queue_parallel_parent_history(
175
- agent: "Agent",
176
- *,
177
- content: Any,
178
- tokens: int = 0,
179
- ) -> bool:
180
- if not is_parallel_worker(agent):
181
- return False
182
- context = agent.context
183
- parent_context_id = str(
184
- context.get_data(PARALLEL_WORKER_PARENT_CONTEXT_KEY) or ""
185
- )
186
- job_id = str(context.get_data(PARALLEL_WORKER_JOB_KEY) or "")
187
- job = _get_job(parent_context_id, job_id)
188
- if not job or job.kind != "tool":
189
- return False
190
- job.parent_history.append((content, tokens))
191
- return True
192
-
193
-
194
- def _jobs_for_context(context: "AgentContext") -> dict[str, ParallelJob]:
195
- jobs = context.get_data(PARALLEL_JOBS_KEY)
196
- if not isinstance(jobs, dict):
197
- jobs = {}
198
- context.set_data(PARALLEL_JOBS_KEY, jobs)
199
- return jobs
200
-
201
-
202
- def _get_job(parent_context_id: str, job_id: str) -> ParallelJob | None:
203
- from agent import AgentContext
204
-
205
- context = AgentContext.get(parent_context_id)
206
- if not context:
207
- return None
208
- job = _jobs_for_context(context).get(job_id)
209
- return job if isinstance(job, ParallelJob) else None
210
-
211
-
212
- def _new_job_id(tool_name: str) -> str:
213
- prefix = "".join(ch for ch in tool_name if ch.isalnum())[:12] or "job"
214
- return f"{prefix}-{uuid.uuid4().hex[:8]}"
215
-
216
-
217
- async def start_parallel_jobs(
218
- agent: "Agent",
219
- calls: list[NormalizedToolCall],
220
- ) -> list[ParallelJob]:
221
- jobs: list[ParallelJob] = []
222
- context = agent.context
223
- job_store = _jobs_for_context(context)
224
-
225
- for call in calls:
226
- kind: JobKind = "subordinate" if call.tool_name == "call_subordinate" else "tool"
227
- job = ParallelJob(
228
- id=_new_job_id(call.tool_name),
229
- parent_context_id=context.id,
230
- index=call.index,
231
- tool_name=call.tool_name,
232
- tool_args=call.tool_args,
233
- kind=kind,
234
- parent_agent=agent,
235
- )
236
- job_store[job.id] = job
237
- jobs.append(job)
238
- _log_parallel_child_started(agent, job)
239
-
240
- try:
241
- job.state = "running"
242
- job.started_at = time.time()
243
- task = DeferredTask(thread_name=THREAD_BACKGROUND)
244
- job.deferred_task = task
245
- if _parallel_worker_kind(agent) == "subordinate" and context.task:
246
- context.task.add_child_task(task)
247
- task.start_task(_run_parallel_job, context.id, job.id)
248
- except Exception as exc:
249
- _finish_job(job, "error", error=str(exc))
250
-
251
- return jobs
252
-
253
-
254
- async def await_parallel_jobs(
255
- agent: "Agent",
256
- job_ids: list[str],
257
- timeout: int = DEFAULT_TIMEOUT_SECONDS,
258
- *,
259
- collect: bool = True,
260
- wait: bool = True,
261
- ) -> list[dict[str, Any]]:
262
- if not job_ids:
263
- raise ValueError("No `job_ids` were provided to await.")
264
-
265
- deadline = time.time() + timeout
266
- wait_timed_out_job_ids: set[str] = set()
267
- while True:
268
- await refresh_parallel_jobs(agent)
269
- jobs = [_jobs_for_context(agent.context).get(job_id) for job_id in job_ids]
270
- missing = [job_id for job_id, job in zip(job_ids, jobs) if job is None]
271
- if missing:
272
- raise ValueError(f"Unknown parallel job id(s): {', '.join(missing)}")
273
-
274
- active = [job for job in jobs if job and job.state not in TERMINAL_STATES]
275
- if not wait or not active:
276
- break
277
-
278
- if time.time() >= deadline:
279
- wait_timed_out_job_ids = {job.id for job in active}
280
- break
281
-
282
- await asyncio.sleep(POLL_INTERVAL_SECONDS)
283
-
284
- snapshots = []
285
- for job_id in job_ids:
286
- job = _jobs_for_context(agent.context).get(job_id)
287
- if job:
288
- snapshot = _job_snapshot(job, include_result=True)
289
- if job.id in wait_timed_out_job_ids and job.state not in TERMINAL_STATES:
290
- snapshot["wait_timed_out"] = True
291
- snapshots.append(snapshot)
292
-
293
- if collect:
294
- await collect_parallel_jobs(agent, job_ids)
295
-
296
- return snapshots
297
-
298
-
299
- async def cancel_parallel_jobs(agent: "Agent", job_ids: list[str]) -> list[dict[str, Any]]:
300
- if not job_ids:
301
- raise ValueError("No `job_ids` were provided to cancel.")
302
-
303
- await refresh_parallel_jobs(agent)
304
- snapshots = []
305
- for job_id in job_ids:
306
- job = _jobs_for_context(agent.context).get(job_id)
307
- if not job:
308
- raise ValueError(f"Unknown parallel job id: {job_id}")
309
- await _cancel_job(job)
310
- snapshots.append(_job_snapshot(job, include_result=True))
311
- await cleanup_parallel_job(agent, job)
312
- _jobs_for_context(agent.context).pop(job_id, None)
313
- return snapshots
314
-
315
-
316
- async def refresh_parallel_jobs(agent: "Agent") -> list[ParallelJob]:
317
- jobs = list(_jobs_for_context(agent.context).values())
318
- for job in jobs:
319
- if job.state in TERMINAL_STATES:
320
- continue
321
- task = job.deferred_task
322
- if not task:
323
- continue
324
- if task.is_ready():
325
- try:
326
- await task.result()
327
- except asyncio.CancelledError:
328
- _finish_job(job, "cancelled", error="Parallel job was cancelled.")
329
- except Exception as exc:
330
- _finish_job(job, "error", error=str(exc))
331
- elif task.is_alive():
332
- job.state = "running"
333
- if job.started_at is None:
334
- job.started_at = time.time()
335
- return jobs
336
-
337
-
338
- async def cleanup_parallel_job(agent: "Agent", job: ParallelJob) -> None:
339
- if job.deferred_task and job.deferred_task.is_alive():
340
- job.deferred_task.kill()
341
- if job.kind == "tool":
342
- await _remove_context(job.worker_context_id)
343
-
344
-
345
- async def collect_parallel_jobs(
346
- agent: "Agent",
347
- job_ids: list[str],
348
- *,
349
- promote_parent_history: bool = False,
350
- ) -> None:
351
- jobs = _jobs_for_context(agent.context)
352
- for job_id in dict.fromkeys(job_ids):
353
- job = jobs.get(job_id)
354
- if not job or job.state not in TERMINAL_STATES:
355
- continue
356
- if promote_parent_history:
357
- for content, tokens in job.parent_history:
358
- agent.hist_add_message(False, content=content, tokens=tokens)
359
- job.parent_history.clear()
360
- await cleanup_parallel_job(agent, job)
361
- jobs.pop(job_id, None)
362
-
363
-
364
- async def build_parallel_jobs_extras(agent: "Agent") -> str:
365
- await refresh_parallel_jobs(agent)
366
- jobs = [
367
- job
368
- for job in _jobs_for_context(agent.context).values()
369
- if isinstance(job, ParallelJob) and job.state not in {"cancelled", "timeout"}
370
- ]
371
- if not jobs:
372
- return ""
373
-
374
- active = [job for job in jobs if job.state not in TERMINAL_STATES]
375
- ready = [job for job in jobs if job.state in TERMINAL_STATES]
376
- if not active and not ready:
377
- return ""
378
-
379
- lines = ["parallel jobs:"]
380
- if active:
381
- lines.append("running:")
382
- for job in active:
383
- lines.append(
384
- f"- {job.id}: {job.tool_name} [{job.state}], running for {job.elapsed():.1f}s"
385
- )
386
- if ready:
387
- lines.append("ready to collect with `parallel` and `job_ids`:")
388
- for job in ready:
389
- lines.append(
390
- f"- {job.id}: {job.tool_name} [{job.state}], duration {job.elapsed():.1f}s"
391
- )
392
- lines.append("call the `parallel` tool with `job_ids` to await/collect results or `action: \"cancel\"` to cancel.")
393
- return "\n".join(lines)
394
-
395
-
396
- def format_started_jobs(jobs: list[ParallelJob]) -> str:
397
- payload = {
398
- "status": "started",
399
- "jobs": [_job_snapshot(job, include_result=False) for job in jobs],
400
- "instruction": "Use the parallel tool with job_ids to await or cancel these background jobs.",
401
- }
402
- return json.dumps(payload, ensure_ascii=False, separators=(",", ":"))
403
-
404
-
405
- def format_parallel_results(results: list[dict[str, Any]]) -> str:
406
- states = [result.get("state") for result in results]
407
- has_active_jobs = any(state not in TERMINAL_STATES for state in states)
408
- wait_timed_out = any(result.get("wait_timed_out") for result in results)
409
- if has_active_jobs:
410
- status = "waiting" if wait_timed_out else "running"
411
- elif states and all(state == "success" for state in states):
412
- status = "success"
413
- elif states and all(state == "cancelled" for state in states):
414
- status = "cancelled"
415
- elif any(state == "success" for state in states):
416
- status = "partial"
417
- else:
418
- status = "error"
419
-
420
- payload = {
421
- "status": status,
422
- "jobs": results,
423
- }
424
- if wait_timed_out:
425
- payload["wait_timeout"] = True
426
- if has_active_jobs:
427
- payload["instruction"] = (
428
- "Some jobs are still running. Call `parallel` with `action: \"await\"` "
429
- "and the listed `job_ids` to wait again, or `action: \"cancel\"` to stop them."
430
- )
431
- return json.dumps(payload, ensure_ascii=False, separators=(",", ":"))
432
-
433
-
434
- async def _run_parallel_job(parent_context_id: str, job_id: str) -> None:
435
- job = _get_job(parent_context_id, job_id)
436
- if not job:
437
- return
438
- try:
439
- if job.kind == "subordinate":
440
- result = await _run_subordinate_context_job(parent_context_id, job)
441
- else:
442
- result = await _run_direct_tool_job(parent_context_id, job)
443
- _finish_job(job, "success", result=result)
444
- except asyncio.CancelledError:
445
- _finish_job(job, "cancelled", error="Parallel job was cancelled.")
446
- raise
447
- except Exception as exc:
448
- _finish_job(job, "error", error=str(exc))
449
- PrintStyle.error(f"Parallel job {job.id} failed: {exc}")
450
-
451
-
452
- async def _run_subordinate_context_job(parent_context_id: str, job: ParallelJob) -> str:
453
- from agent import AgentContext
454
- from helpers.tool_policy import ensure_tool_allowed
455
- from tools.call_subordinate import get_or_create_subordinate, run_subordinate
456
-
457
- parent_context = AgentContext.get(parent_context_id)
458
- if not parent_context:
459
- raise ValueError("Parent context not found.")
460
- parent_agent = job.parent_agent or parent_context.agent0
461
- ensure_tool_allowed(parent_agent, "call_subordinate")
462
-
463
- args = job.tool_args
464
- message = str(args.get("message") or "").strip()
465
- if not message:
466
- raise ValueError("call_subordinate requires `tool_args.message`.")
467
-
468
- context_id = str(args.get("context_id") or args.get("agent_id") or "").strip()
469
- reset = args.get("reset", False)
470
- slot = (
471
- job.id
472
- if coerce_bool(reset, False) and not context_id
473
- else "default"
474
- )
475
- attachments = args.get("attachments") if isinstance(args.get("attachments"), list) else []
476
- subordinate = get_or_create_subordinate(
477
- parent_agent,
478
- profile=str(args.get("profile") or args.get("agent_profile") or ""),
479
- reset=reset,
480
- context_id=context_id,
481
- name=str(args.get("name") or ""),
482
- message=message,
483
- slot=slot,
484
- )
485
- worker_context = subordinate.context
486
- job.worker_context_id = worker_context.id
487
- if job.deferred_task:
488
- worker_context.task = job.deferred_task
489
-
490
- worker_context.set_data(PARALLEL_WORKER_PARENT_CONTEXT_KEY, parent_context.id)
491
- worker_context.set_data(PARALLEL_WORKER_JOB_KEY, job.id)
492
- worker_context.set_data(PARALLEL_WORKER_KIND_KEY, job.kind)
493
- worker_context.set_output_data(CHILD_PARALLEL_JOB_ID_KEY, job.id)
494
- worker_context.set_output_data(CHILD_PARALLEL_TOOL_NAME_KEY, job.tool_name)
495
- return await run_subordinate(parent_agent, subordinate, message, attachments)
496
-
497
-
498
- async def _run_direct_tool_job(parent_context_id: str, job: ParallelJob) -> str:
499
- from agent import AgentContext, AgentContextType, LoopData
500
-
501
- parent_context = AgentContext.get(parent_context_id)
502
- if not parent_context:
503
- raise ValueError("Parent context not found.")
504
-
505
- worker_context: AgentContext | None = None
506
- try:
507
- worker_context = AgentContext(
508
- config=_clone_config(parent_context.config),
509
- name=f"parallel:{job.tool_name}",
510
- type=AgentContextType.BACKGROUND,
511
- )
512
- worker_context.set_data(PARALLEL_WORKER_PARENT_CONTEXT_KEY, parent_context_id)
513
- worker_context.set_data(PARALLEL_WORKER_JOB_KEY, job.id)
514
- worker_context.set_data(PARALLEL_WORKER_KIND_KEY, job.kind)
515
- worker_context.set_data(
516
- "chat_model_override",
517
- parent_context.get_data("chat_model_override"),
518
- )
519
- job.worker_context_id = worker_context.id
520
- _copy_project(parent_context, worker_context)
521
-
522
- worker_agent = worker_context.agent0
523
- worker_agent.last_user_message = parent_context.agent0.last_user_message
524
- worker_agent.loop_data = LoopData()
525
- return await execute_tool_call(
526
- worker_agent,
527
- job.tool_name,
528
- job.tool_args,
529
- log_item=job.log_item,
530
- )
531
- finally:
532
- if worker_context:
533
- await _remove_context(worker_context.id)
534
-
535
-
536
- async def execute_tool_call(
537
- agent: "Agent",
538
- tool_name: str,
539
- tool_args: dict[str, Any],
540
- *,
541
- log_item: "LogItem | None" = None,
542
- ) -> str:
543
- if tool_name == "parallel":
544
- raise ValueError("`parallel` cannot be nested inside a parallel worker.")
545
-
546
- tool = _resolve_parallel_tool(agent, tool_name, tool_args, strict=True)
547
- if not tool:
548
- raise ValueError(f"Tool '{tool_name}' not found or could not be initialized.")
549
-
550
- original_get_log_object = None
551
- if log_item is not None:
552
- original_get_log_object = tool.get_log_object
553
- tool.get_log_object = lambda: log_item
554
-
555
- agent.loop_data.current_tool = tool
556
- try:
557
- await agent.handle_intervention()
558
- await tool.before_execution(**tool_args)
559
- await agent.handle_intervention()
560
- await call_extensions_async(
561
- "tool_execute_before",
562
- agent,
563
- tool_args=tool_args or {},
564
- tool_name=tool_name,
565
- )
566
- response = await tool.execute(**tool_args)
567
- await agent.handle_intervention()
568
- await call_extensions_async(
569
- "tool_execute_after",
570
- agent,
571
- response=response,
572
- tool_name=tool_name,
573
- )
574
- await tool.after_execution(response)
575
- await agent.handle_intervention()
576
- return response.message
577
- finally:
578
- if original_get_log_object is not None:
579
- tool.get_log_object = original_get_log_object
580
- agent.loop_data.current_tool = None
581
-
582
-
583
- def _resolve_parallel_tool(
584
- agent: "Agent",
585
- tool_name: str,
586
- tool_args: dict[str, Any],
587
- *,
588
- strict: bool = False,
589
- ):
590
- message = json.dumps({"tool_name": tool_name, "tool_args": tool_args})
591
-
592
- tool = None
593
- try:
594
- import helpers.mcp_handler as mcp_helper
595
-
596
- tool = mcp_helper.MCPConfig.get_instance().get_tool(agent, tool_name)
597
- except ImportError:
598
- tool = None
599
- except Exception as exc:
600
- if strict:
601
- raise
602
- PrintStyle.warning(f"Failed to initialize MCP tool '{tool_name}' for parallel job: {exc}")
603
-
604
- if not tool:
605
- get_tool = getattr(agent, "get_tool", None)
606
- if not callable(get_tool):
607
- if strict:
608
- raise ValueError(f"Tool '{tool_name}' not found or could not be initialized.")
609
- return None
610
- try:
611
- tool = get_tool(
612
- name=tool_name,
613
- method=None,
614
- args=tool_args,
615
- message=message,
616
- loop_data=getattr(agent, "loop_data", None),
617
- )
618
- except Exception as exc:
619
- if strict:
620
- raise
621
- PrintStyle.warning(f"Failed to initialize tool '{tool_name}' for parallel job: {exc}")
622
- tool = None
623
-
624
- if not tool:
625
- return None
626
-
627
- try:
628
- tool.args = dict(tool_args)
629
- except Exception:
630
- pass
631
-
632
- return tool
633
-
634
-
635
- async def _cancel_job(
636
- job: ParallelJob,
637
- *,
638
- state: JobState = "cancelled",
639
- message: str = "Parallel job was cancelled.",
640
- ) -> None:
641
- if job.deferred_task and job.deferred_task.is_alive():
642
- job.deferred_task.kill()
643
- _finish_job(job, state, error=message)
644
-
645
-
646
- def _finish_job(
647
- job: ParallelJob,
648
- state: JobState,
649
- *,
650
- result: str | None = None,
651
- error: str | None = None,
652
- ) -> None:
653
- job.state = state
654
- job.completed_at = time.time()
655
- if job.started_at is None:
656
- job.started_at = job.created_at
657
- if result is not None:
658
- job.result = result
659
- if error is not None:
660
- job.error = error
661
- _update_parallel_child_log(job)
662
-
663
-
664
- async def _remove_context(context_id: str | None) -> None:
665
- if not context_id:
666
- return
667
- from agent import AgentContext
668
- from helpers import persist_chat
669
-
670
- context = AgentContext.get(context_id)
671
- if context:
672
- try:
673
- context.reset()
674
- except Exception:
675
- pass
676
- AgentContext.remove(context_id)
677
- persist_chat.remove_chat(context_id)
678
-
679
-
680
- def _log_parallel_child_started(agent: "Agent", job: ParallelJob) -> None:
681
- if job.kind == "subordinate":
682
- job.log_item = agent.context.log.log(
683
- type="subagent",
684
- heading=f"icon://communication {agent.agent_name}: Calling Subordinate Agent",
685
- content="",
686
- kvps=job.tool_args,
687
- id=job.log_id,
688
- )
689
- return
690
-
691
- tool = _resolve_parallel_tool(agent, job.tool_name, job.tool_args)
692
- if tool is not None:
693
- try:
694
- job.log_item = tool.get_log_object()
695
- if job.log_item is not None:
696
- return
697
- except Exception as exc:
698
- PrintStyle.warning(
699
- f"Failed to derive parallel child log for {job.tool_name}: {exc}"
700
- )
701
-
702
- heading = f"icon://construction {agent.agent_name}: Using tool '{job.tool_name}'"
703
- job.log_item = agent.context.log.log(
704
- type="tool",
705
- heading=heading,
706
- content="",
707
- kvps=job.tool_args,
708
- id=job.log_id,
709
- _tool_name=job.tool_name,
710
- )
711
-
712
-
713
- def _update_parallel_child_log(job: ParallelJob) -> None:
714
- if not job.log_item:
715
- return
716
- if job.state == "success":
717
- content = job.result if job.result else "(completed without textual output)"
718
- else:
719
- content = f"Error: {job.error or job.state}"
720
- try:
721
- job.log_item.update(content=content)
722
- except Exception:
723
- pass
724
-
725
-
726
- def _job_snapshot(job: ParallelJob, *, include_result: bool) -> dict[str, Any]:
727
- data: dict[str, Any] = {
728
- "job_id": job.id,
729
- "tool_name": job.tool_name,
730
- "state": job.state,
731
- "duration_seconds": round(job.elapsed(), 3),
732
- }
733
- if job.worker_context_id:
734
- data["context_id"] = job.worker_context_id
735
- if include_result:
736
- if job.result is not None:
737
- data["result"] = job.result
738
- if job.error is not None:
739
- data["error"] = job.error
740
- return data
741
-
742
-
743
- def _clone_config(config: "AgentConfig") -> "AgentConfig":
744
- try:
745
- return replace(
746
- config,
747
- knowledge_subdirs=list(config.knowledge_subdirs),
748
- additional=dict(config.additional),
749
- )
750
- except Exception:
751
- return config
752
-
753
-
754
- def _copy_project(parent_context: "AgentContext", worker_context: "AgentContext") -> None:
755
- try:
756
- from helpers import projects
757
-
758
- project_name = projects.get_context_project_name(parent_context)
759
- if project_name:
760
- projects.activate_project(worker_context.id, project_name, mark_dirty=False)
761
- except Exception:
762
- pass
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:69341a3cf61a8a259e140e06cdf13248d4a25b42c6a9cd9709debdf9d23f69ea
3
+ size 25127
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/parallel_tools.py.dox.md CHANGED
@@ -1,66 +1,3 @@
1
- # parallel_tools.py DOX
2
-
3
- ## Purpose
4
-
5
- - Own the shared runtime for parallel tool-call jobs.
6
- - Normalize wrapped tool-call payloads, start background jobs, await or cancel jobs, and render prompt extras for active parallel work.
7
- - Keep this file-level DOX profile synchronized with `parallel_tools.py` because this directory is intentionally flat.
8
-
9
- ## Ownership
10
-
11
- - `parallel_tools.py` owns the runtime implementation.
12
- - `parallel_tools.py.dox.md` owns durable notes about responsibilities, contracts, side effects, and verification for that implementation.
13
- - Public concepts:
14
- - `NormalizedToolCall`
15
- - `ParallelJob`
16
- - `start_parallel_jobs(...)`
17
- - `await_parallel_jobs(...)`
18
- - `cancel_parallel_jobs(...)`
19
- - `build_parallel_jobs_extras(...)`
20
- - `format_parallel_results(...)`
21
-
22
- ## Runtime Contracts
23
-
24
- - Helper modules own reusable framework APIs and must preserve public callers unless all callers, tests, and docs are updated together.
25
- - Wrapped tool-call items must use the same shape as normal tool calls: a tool name plus arguments.
26
- - Normalization accepts full agent-reply-shaped objects when `tool_name` and `tool_args` are present; non-contract planning fields such as `thoughts` or `headline` are ignored.
27
- - `tool_calls` should be an array, but normalization also accepts a valid JSON string encoding of that array to recover provider/model stringification.
28
- - Normalization rejects `document_query` and `response` inside `parallel`: document parsing and Q&A must run sequentially, while `response` must remain top-level so it can end the message loop.
29
- - `call_subordinate` jobs first enforce the actual calling agent's delegation policy, then call the same creation and execution functions as direct delegation in `tools/call_subordinate.py`; this helper does not construct or prompt a second kind of subordinate.
30
- - Fresh parallel sibling calls create distinct `parent.number + 1` child agents. Their job snapshots expose stable `context_id` values that direct or parallel `reset=false` calls can continue after success or failure.
31
- - Jobs retain their actual parent agent so parallel calls made by A1 create A2 rather than falling back to a context's A0.
32
- - Subordinate child chats are tagged with job metadata, remain outside the scheduler task list, and may use normal child-chat tools including `parallel`.
33
- - Nested parallel jobs started by a parallel subordinate are registered as child `DeferredTask` instances so stopping the ancestor also stops its descendants.
34
- - Direct tool jobs run in isolated background contexts and are blocked from recursively invoking `parallel`.
35
- - Direct tool jobs inherit the parent's active per-chat model override and current user message.
36
- - Direct tool background context cleanup removes both the in-memory context and any transient chat folder left on disk.
37
- - Parent-visible child log items are created for each wrapped call so the WebUI can inspect concurrent children separately while the wrapper result remains model-history-only.
38
- - Child tool logs mirror normal tool-call visible args; job ids remain available through wrapper results and prompt extras rather than visible process-step args.
39
- - Wrapped tool child logs use each tool's native `get_log_object()` output when available, preserving special log rendering (for example: `code_execution_tool` uses `code_exe`, `wait` uses `progress`, MCP tools use `mcp`, and regular tools use `tool`).
40
- - Direct parallel worker execution reuses the parent-visible child log item so tool `before_execution()` cannot create a second generic worker log or lose the native badge type.
41
- - Direct tools may explicitly queue model-visible history for their parent. Terminal collection records the outer `parallel` result first, promotes queued messages in job order, and only then removes disposable worker state; background jobs retain queued history until they are collected.
42
- - Job IDs are stable handles for later await, collect, or cancel operations.
43
- - Prompt extras must stay bounded and expose only job IDs, tool names, status, and compact result/error summaries.
44
-
45
- ## Key Concepts
46
-
47
- - The parent context stores in-flight jobs under a private data key; collected terminal jobs are removed from that registry.
48
- - `wait=True` starts jobs and awaits them before returning until all requested jobs finish or the wait timeout is reached; the timeout stops waiting but does not cancel running jobs.
49
- - `collect` returns already-finished job results without waiting; `await` waits for requested job IDs.
50
- - Canceled jobs should be marked terminal and should stop their background `DeferredTask` when cancellation is possible.
51
- - `queue_parallel_parent_history(...)` accepts messages only from registered direct tool workers. `collect_parallel_jobs(...)` optionally promotes those messages while collecting terminal jobs; arbitrary worker history is never copied.
52
-
53
- ## Work Guidance
54
-
55
- - Keep normalization compatible with provider tool-call envelopes and direct JSON objects.
56
- - Avoid importing heavy runtime modules at import time unless startup behavior is verified.
57
- - Coordinate argument, output, or status changes with `tools/parallel.py`, prompt instructions, and tests.
58
-
59
- ## Verification
60
-
61
- - Run targeted tests for normalization, recursion guard, prompt extras, and tool result formatting.
62
- - Run a live Agent Zero chat when changing parallel execution, child chat metadata, or subordinate task behavior.
63
-
64
- ## Child DOX Index
65
-
66
- No child DOX files.
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:3a103b1cbfb2e49767793916f631bee0f5d3923189250a3de2158229e906d21e
3
+ size 5412
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/performance.py CHANGED
@@ -1,51 +1,3 @@
1
- import functools
2
- import inspect
3
- from pyinstrument import Profiler
4
-
5
- def trace_performance(*, show_all=False, color=True, unicode=True):
6
- """
7
- Decorator that profiles a function and prints a call tree when it finishes.
8
-
9
- Works with both synchronous and asynchronous functions.
10
- """
11
-
12
- def decorator(func):
13
- is_coro = inspect.iscoroutinefunction(func)
14
-
15
- @functools.wraps(func)
16
- async def async_wrapper(*args, **kwargs):
17
- profiler = Profiler()
18
- profiler.start()
19
- try:
20
- return await func(*args, **kwargs)
21
- finally:
22
- profiler.stop()
23
- print(f"\n=== Performance trace: {func.__module__}.{func.__qualname__} (async) ===")
24
- print(
25
- profiler.output_text(
26
- color=color,
27
- unicode=unicode,
28
- show_all=show_all,
29
- )
30
- )
31
-
32
- @functools.wraps(func)
33
- def sync_wrapper(*args, **kwargs):
34
- profiler = Profiler()
35
- profiler.start()
36
- try:
37
- return func(*args, **kwargs)
38
- finally:
39
- profiler.stop()
40
- print(f"\n=== Performance trace: {func.__module__}.{func.__qualname__} ===")
41
- print(
42
- profiler.output_text(
43
- color=color,
44
- unicode=unicode,
45
- show_all=show_all,
46
- )
47
- )
48
-
49
- return async_wrapper if is_coro else sync_wrapper
50
-
51
- return decorator
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:198b8cebbf85bdc0b7a6068a4727f14a2efae02c3475e73f069dd17b143f9c9d
3
+ size 1615
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
helpers/performance.py.dox.md CHANGED
@@ -1,42 +1,3 @@
1
- # performance.py DOX
2
-
3
- ## Purpose
4
-
5
- - Own the `performance.py` helper module.
6
- - This module provides lightweight performance tracing helpers.
7
- - Keep this file-level DOX profile synchronized with `performance.py` because this directory is intentionally flat.
8
-
9
- ## Ownership
10
-
11
- - `performance.py` owns the runtime implementation.
12
- - `performance.py.dox.md` owns durable notes about responsibilities, contracts, side effects, and verification for that implementation.
13
- - Top-level functions:
14
- - `trace_performance(show_all=..., color=..., unicode=...)`: Decorator that profiles a function and prints a call tree when it finishes.
15
-
16
- ## Runtime Contracts
17
-
18
- - Helper modules own reusable framework APIs and must preserve public callers unless all callers, tests, and docs are updated together.
19
- - Update this file whenever public functions, classes, persistence behavior, path/security assumptions, side effects, or cross-module contracts change.
20
- - Imported dependency areas include: `functools`, `inspect`, `pyinstrument`.
21
-
22
- ## Key Concepts
23
-
24
- - Important called helpers/classes observed in the source: `inspect.iscoroutinefunction`, `functools.wraps`, `Profiler`, `profiler.start`, `profiler.stop`, `func`, `profiler.output_text`.
25
- - Keep request/response, tool, or helper semantics documented here at the same time as source changes.
26
-
27
- ## Work Guidance
28
-
29
- - Preserve public helper APIs used by core code and plugins unless every caller is updated.
30
- - Keep path, auth, secret, persistence, network, and subprocess behavior explicit and bounded.
31
- - Prefer adding cohesive helper functions here only when behavior is reused across modules.
32
-
33
- ## Verification
34
-
35
- - Run targeted tests for changed helper behavior; run security regressions for auth, filesystem, WebSocket, tunnel, upload, or secret-handling helpers.
36
- - Related tests observed by source search:
37
- - `tests/test_extensions_stress.py`
38
- - `tests/test_ws_manager.py`
39
-
40
- ## Child DOX Index
41
-
42
- No child DOX files.
 
1
+ version https://git-lfs.github.com/spec/v1
2
+ oid sha256:81afd1008693788b59b3d810ded2379dc682f5cee0d74a9bd30d7959774f8c15
3
+ size 1937