-
Notifications
You must be signed in to change notification settings - Fork 36
feat: add streaming output for dp train #307
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| @@ -1,4 +1,7 @@ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| import os | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| import subprocess | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| import sys | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| import threading | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| from typing import ( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| List, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Tuple, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -11,6 +14,74 @@ | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| from dflow.utils import run_command as dflow_run_command | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| def run_command_streaming( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| cmd: Union[str, List[str]], | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| shell: bool = False, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| log_file=None, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ) -> Tuple[int, str, str]: | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| """Run command with streaming output to both terminal and log file.""" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if isinstance(cmd, str): | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| cmd = cmd if shell else cmd.split() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
Comment on lines
+23
to
+25
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🛠️ Refactor suggestion Use shlex for robust splitting and handle shell=True with list commands
+import shlex
@@
- if isinstance(cmd, str):
- cmd = cmd if shell else cmd.split()
+ if isinstance(cmd, str):
+ cmd = cmd if shell else shlex.split(cmd)
+ elif shell and isinstance(cmd, list):
+ # When shell=True, pass a single string; preserve quoting.
+ cmd = shlex.join(cmd)📝 Committable suggestion
Suggested change
🤖 Prompt for AI Agents |
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| # Open log file if specified | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| log_fp = open(log_file, "w") if log_file else None | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| try: | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| # Start subprocess | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| process = subprocess.Popen( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| cmd, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| stdout=subprocess.PIPE, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| stderr=subprocess.PIPE, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| stdin=subprocess.PIPE, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| shell=shell, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| text=True, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| bufsize=1, # Line buffered | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| universal_newlines=True, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| # Store output | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| stdout_buffer = [] | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| stderr_buffer = [] | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| def stream_output(pipe, buffer, is_stderr=False): | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| for line in iter(pipe.readline, ""): | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| buffer.append(line) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. P1 — Bound the in-memory capture for multi-day training. Every stdout/stderr line is appended here until the process exits, and line 78 then joins each complete buffer into a second full-size string. The new helper therefore keeps memory proportional to the entire multi-day log and can OOM the wrapper even though the complete stream is already on disk. Keep only a bounded tail (for example, deques used for failure diagnostics) and stream the full output exclusively to the log. This needs an import plus changes to the return contract, so a local suggestion would not be complete. Codex quota is about to reset, so I am using the remaining token budget to review this PR now. Coding agent: Codex |
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| # Print to terminal | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if is_stderr: | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| print(line, end="", file=sys.stderr) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| else: | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| print(line, end="") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| # Write to log file | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if log_fp: | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| log_fp.write(line) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| log_fp.flush() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| pipe.close() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
Comment on lines
+42
to
+58
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Guard concurrent writes to the same log file
# Store output
stdout_buffer = []
stderr_buffer = []
+ log_lock = threading.Lock() if log_fp else None
def stream_output(pipe, buffer, is_stderr=False):
for line in iter(pipe.readline, ""):
buffer.append(line)
# Print to terminal
if is_stderr:
print(line, end="", file=sys.stderr)
else:
print(line, end="")
# Write to log file
if log_fp:
- log_fp.write(line)
- log_fp.flush()
+ with log_lock:
+ log_fp.write(line)
+ log_fp.flush()
pipe.close()📝 Committable suggestion
Suggested change
🤖 Prompt for AI Agents |
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| # Start threads for streaming | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| stdout_thread = threading.Thread( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| target=stream_output, args=(process.stdout, stdout_buffer, False) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| stderr_thread = threading.Thread( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| target=stream_output, args=(process.stderr, stderr_buffer, True) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| stdout_thread.start() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| stderr_thread.start() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| # Wait for process to complete | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return_code = process.wait() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| # Wait for threads to finish | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| stdout_thread.join() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| stderr_thread.join() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return return_code, "".join(stdout_buffer), "".join(stderr_buffer) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| finally: | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if log_fp: | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| log_fp.close() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| def run_command( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| cmd: Union[str, List[str]], | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| shell: bool = False, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
P1 — Use a single owner for
train.log. The file is already open asfplogat line 270.run_command_streamingopens the same path again with"w"and writes the training output, but the original descriptor still has offset 0; the later TensorFlow freeze writes at lines 346-349 therefore overwrite the beginning of the streamed training log. Remove/close the original handle before streaming and append freeze output afterward, or route both phases through one logging owner. The fix spans the unchanged open plus both phases, so a one-line suggestion would be incomplete.Codex quota is about to reset, so I am using the remaining token budget to review this PR now.
Coding agent: Codex
Codex version: codex-cli 0.144.6
Model: gpt-5.6-sol
Reasoning effort: xhigh