Throttle and schedule buffered output delivery
This commit is contained in:
@@ -30,6 +30,7 @@ class Bridge:
|
||||
self._last_sent_fingerprint: str | None = None
|
||||
self._last_input_at = time.monotonic()
|
||||
self._last_output_at = time.monotonic()
|
||||
self._output_buffer_started_at: float | None = None
|
||||
self._input_idle_reported = False
|
||||
self._output_idle_reported = False
|
||||
|
||||
@@ -118,6 +119,7 @@ class Bridge:
|
||||
self._running = True
|
||||
self._last_input_at = time.monotonic()
|
||||
self._last_output_at = time.monotonic()
|
||||
self._output_buffer_started_at = None
|
||||
self._input_idle_reported = False
|
||||
self._output_idle_reported = False
|
||||
self._output_thread = threading.Thread(target=self._poll_output, daemon=True)
|
||||
@@ -137,6 +139,7 @@ class Bridge:
|
||||
self._active_target = None
|
||||
self._last_sent_output = ""
|
||||
self._last_sent_fingerprint = None
|
||||
self._output_buffer_started_at = None
|
||||
self.slack.send_message(channel, ":electric_plug: 세션 연결이 해제되었습니다.")
|
||||
|
||||
def _poll_output(self) -> None:
|
||||
@@ -144,17 +147,27 @@ class Bridge:
|
||||
buffer = ""
|
||||
while self._running and self.pty and self.pty.is_alive:
|
||||
now = time.monotonic()
|
||||
output = self.pty.read_output(timeout=self.config.pty_read_timeout)
|
||||
if buffer and self._should_flush_output_buffer(now):
|
||||
self._send_output_chunks(buffer)
|
||||
buffer = ""
|
||||
self._output_buffer_started_at = None
|
||||
|
||||
read_timeout = self._next_read_timeout(now, has_buffer=bool(buffer))
|
||||
output = self.pty.read_output(timeout=read_timeout)
|
||||
now = time.monotonic()
|
||||
if output:
|
||||
cleaned = clean_terminal_output(output)
|
||||
if cleaned:
|
||||
buffer += f"{cleaned}\n"
|
||||
self._last_output_at = now
|
||||
if self._output_buffer_started_at is None:
|
||||
self._output_buffer_started_at = now
|
||||
self._output_idle_reported = False
|
||||
|
||||
if buffer:
|
||||
if buffer and self._should_flush_output_buffer(now):
|
||||
self._send_output_chunks(buffer)
|
||||
buffer = ""
|
||||
self._output_buffer_started_at = None
|
||||
|
||||
output_idle = now - self._last_output_at
|
||||
if (
|
||||
@@ -186,11 +199,51 @@ class Bridge:
|
||||
time.sleep(self.config.output_buffer_interval)
|
||||
|
||||
if not self._running:
|
||||
self._output_buffer_started_at = None
|
||||
return
|
||||
|
||||
if buffer:
|
||||
self._send_output_chunks(buffer)
|
||||
self._output_buffer_started_at = None
|
||||
|
||||
# attach 프로세스가 예기치 않게 종료된 경우
|
||||
self.slack.send_message(self._channel, ":warning: 세션 연결이 종료되었습니다.")
|
||||
|
||||
def _next_read_timeout(self, now: float, has_buffer: bool) -> float:
|
||||
"""다음 PTY 읽기 타임아웃을 계산한다."""
|
||||
base_timeout = max(0.0, float(self.config.pty_read_timeout))
|
||||
if not has_buffer:
|
||||
return base_timeout
|
||||
|
||||
deadline = self._next_output_flush_deadline()
|
||||
if deadline is None:
|
||||
return base_timeout
|
||||
|
||||
remaining = max(0.0, deadline - now)
|
||||
return min(base_timeout, remaining)
|
||||
|
||||
def _next_output_flush_deadline(self) -> float | None:
|
||||
"""버퍼 flush의 가장 이른 데드라인을 반환한다."""
|
||||
deadlines: list[float] = []
|
||||
settle_seconds = max(0.0, self.config.output_settle_seconds)
|
||||
if settle_seconds > 0:
|
||||
deadlines.append(self._last_output_at + settle_seconds)
|
||||
|
||||
flush_interval_seconds = max(0.0, self.config.output_flush_interval_seconds)
|
||||
if self._output_buffer_started_at is not None:
|
||||
if flush_interval_seconds == 0:
|
||||
return self._output_buffer_started_at
|
||||
deadlines.append(self._output_buffer_started_at + flush_interval_seconds)
|
||||
|
||||
if not deadlines:
|
||||
return None
|
||||
return min(deadlines)
|
||||
|
||||
def _should_flush_output_buffer(self, now: float) -> bool:
|
||||
"""버퍼를 Slack으로 전송할 시점을 계산한다."""
|
||||
deadline = self._next_output_flush_deadline()
|
||||
return deadline is not None and now >= deadline
|
||||
|
||||
@staticmethod
|
||||
def _split_message(text: str, max_length: int) -> list[str]:
|
||||
"""긴 텍스트를 메시지 길이 제한에 맞게 분할한다."""
|
||||
|
||||
@@ -25,6 +25,10 @@ class Config:
|
||||
|
||||
# Buffer
|
||||
output_buffer_interval: float = float(os.getenv("OUTPUT_BUFFER_INTERVAL", "2.0"))
|
||||
output_settle_seconds: float = float(os.getenv("OUTPUT_SETTLE_SECONDS", "4.0"))
|
||||
output_flush_interval_seconds: float = float(
|
||||
os.getenv("OUTPUT_FLUSH_INTERVAL_SECONDS", "15.0")
|
||||
)
|
||||
max_message_length: int = int(os.getenv("MAX_MESSAGE_LENGTH", "3000"))
|
||||
|
||||
# Status reporting / reconnect
|
||||
|
||||
@@ -55,7 +55,7 @@ class PtyManager:
|
||||
logger.debug("입력 전송: %s", text)
|
||||
self._process.sendline(text)
|
||||
|
||||
def read_output(self, timeout: int = 5) -> str:
|
||||
def read_output(self, timeout: float = 5) -> str:
|
||||
"""프로세스의 출력을 읽는다."""
|
||||
if not self.is_alive:
|
||||
raise RuntimeError("프로세스가 실행 중이 아닙니다.")
|
||||
|
||||
Reference in New Issue
Block a user