Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
101 changes: 88 additions & 13 deletions lib/aio/s3streamer.py
Original file line number Diff line number Diff line change
@@ -1,17 +1,92 @@
# Copyright (C) 2022-2024 Red Hat, Inc.
# SPDX-FileCopyrightText: 2022-2024 Red Hat, Inc.
# SPDX-License-Identifier: AGPL-3.0-or-later

# S3 Log Streaming Protocol
# =========================
#
# This module implements a protocol for streaming log output through S3 (or
# any object store with similar semantics). The core problem is that S3
# objects are immutable once written, so we can't append to a file. Instead,
# we write a series of chunk objects and a manifest that tells clients how to
# reassemble them. The base filename ("log" in this case) is not special —
# any name works as long as the writer and clients agree on it.
#
# The entire protocol is bytes-oriented: chunk sizes in the manifest are
# byte counts, Range request offsets are byte offsets, and chunk filenames
# encode byte ranges. Although character offsets are never used for
# anything, the content is assumed to be UTF-8 text. The server

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

here you promise UTF-8 in the entire protocol but

self.destination.write('log.chunks', json.dumps(chunk_sizes).encode('ascii'))

:)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ascii is utf8 :)

# is responsible for ensuring that a single codepoint is never split across
# two chunk files, so clients can safely decode each chunk independently.
#
# The object store should support Range requests and respond with 206
# Partial Content. If the server does not support Range requests (e.g.
# python -m http.server for local development), it will return 200 with
# the complete object body. Clients must handle both cases: on 206 the
# body is exactly the requested range; on 200 the client must slice the
# body at the requested byte offset and discard the prefix.
#
# File layout during streaming:
#
# log.chunks JSON array of chunk sizes in bytes, e.g. [81920, 40960]
# log.0-81920 First chunk (size 81920, bytes [0, 81920))
# log.81920-122880 Second chunk (size 40960, bytes [81920, 122880))
#
# File layout after streaming:
#
# log Complete log contents (single object)
#
# The chunk files and log.chunks are deleted when streaming ends.
#
#
# Writer protocol (this module)
# -----------------------------
#
# LogStreamer buffers incoming data and flushes it as chunk objects. A flush
# happens when the buffer exceeds SIZE_LIMIT (1MB) or TIME_LIMIT (30s) elapses
# with data pending.
#
# Each flush appends a new entry to the internal chunks list, then runs a
# merge pass modelled after the game 2048: if the last two entries have the
# same number of constituent blocks, they are merged into one. This keeps
# the total number of chunk objects logarithmic in the amount of data written,
# while only ever rewriting the last chunk (so the write amplification is
# also logarithmically bounded).
#
# The manifest (log.chunks) is a JSON array of chunk sizes in bytes. Clients
# use it to derive the chunk filenames: each chunk is named
# "log.{start}-{end}" where start is the cumulative size of all preceding
# chunks and end = start + size. The range is end-exclusive: a chunk named
# "log.0-81920" has length 81920, containing bytes [0, 81920).
#
# Order of operations matters for consistency. On each flush, the writer
# must write the chunk object first, then update the manifest — otherwise a
# client could read a manifest that references a chunk that doesn't exist
# yet. On close(), the writer writes the final "log" object first, then
# deletes the manifest and chunk files. This way, a client that gets a 404
# on the manifest can be confident that the final log object is already
# available.
#
#
# Client protocol (log.html, s3stream CLI)
# -----------------------------------------
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU Affero General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
# Clients always read the manifest first. On each poll they compare it
# against how many bytes they have already received, then fetch only the
# new/updated chunk objects, using Range requests to avoid re-downloading
# data within a chunk that grew due to a merge. If the server responds
# with 200 instead of 206, the client reads the full body as bytes and
# slices at the requested offset before decoding to text.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU Affero General Public License for more details.
# S3 can return 500 Internal Server Error at any time, even under normal
# operation. Both clients and the server should retry transient failures (5xx,
# network errors) with exponential backoff:
# https://docs.aws.amazon.com/AmazonS3/latest/developerguide/ErrorBestPractices.html
#
# You should have received a copy of the GNU Affero General Public License
# along with this program. If not, see <http://www.gnu.org/licenses/>.
# When fetching the manifest or a chunk returns 404 (or 403 — S3 returns 403
# instead of 404 when the bucket policy does not grant s3:ListBucket),
# streaming is over (or never started). The client fetches "log" with a Range
# header for any bytes beyond what it already has, giving a seamless transition
# from streamed to final content.

import asyncio
import json
Expand Down Expand Up @@ -137,12 +212,12 @@ def send_pending(self) -> None:

def start(self, data: str) -> None:
# Send the initial data immediately, to get the chunks file written out.
self.pending = data.encode()
self.pending = data.encode(errors='replace')
self.send_pending()
AttachmentsDirectory(self.index, f'{LIB_DIR}/s3-html').scan()

def write(self, data: str) -> None:
self.pending += data.encode()
self.pending += data.encode(errors='replace')

if len(self.pending) > LogStreamer.SIZE_LIMIT:
self.send_pending()
Expand Down
99 changes: 99 additions & 0 deletions s3stream
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
#!/usr/bin/env python3

# SPDX-FileCopyrightText: 2026 Red Hat, Inc.
# SPDX-License-Identifier: GPL-3.0-or-later

import argparse
import contextlib
import json
import logging
import time
import urllib.error
import urllib.request
from collections.abc import Iterator

logger = logging.getLogger(__name__)


def fetch_once(url: str, offset: int) -> str:
req = urllib.request.Request(url, headers={"Range": f"bytes={offset}-"})
with urllib.request.urlopen(req) as resp:
data = resp.read()
if resp.status == 206:
return data.decode()
# Server doesn't support Range (e.g. python -m http.server)
if resp.status == 200:
return data[offset:].decode()
raise ValueError(f"Unexpected status {resp.status}")


def fetch(url: str, offset: int = 0) -> str:
for attempt in range(8):
try:
return fetch_once(url, offset)
except OSError as exc:
if isinstance(exc, urllib.error.HTTPError) and exc.code < 500:
raise
delay = 2 ** attempt
logger.debug("error fetching %r, retrying after %r", url, delay)
time.sleep(delay)

return fetch_once(url, offset)


def fetch_log(base_url: str, follow: bool) -> Iterator[str]:
bytes_received = 0

while True:
try:
chunks: list[int] = json.loads(fetch(f"{base_url}.chunks"))
logger.debug("Chunks: %r, bytes: %r", chunks, bytes_received)

chunk_start = 0
for chunk_size in chunks:
chunk_end = chunk_start + chunk_size
if bytes_received < chunk_end:
yield fetch(
f"{base_url}.{chunk_start}-{chunk_end}",
bytes_received - chunk_start,
)
bytes_received = chunk_end
chunk_start = chunk_end

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This assumes that the fetch call above always returns the full chunk, no? I.e., that bytes_received == chunk_end - chunk_start.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I reworked this a bit, but the general logic here is: we trust that urllib does what we asked it or throws an error. In the event of an incomplete read, we expect to hear about that...

except urllib.error.HTTPError as exc:
# S3 returns 403 instead of 404 when s3:ListBucket is not granted
if exc.code in (403, 404):
break

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does this break the while loop or just the try?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

break doesn't interact with try in any way...

raise

if not follow:
return

time.sleep(10)

# Streaming is over (or never started) — fetch the (remaining) final log

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

So the break above gets us here when .chunks doesn't exist?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

exactly.

yield fetch(base_url, bytes_received)


def main() -> None:
parser = argparse.ArgumentParser(description="Fetch streamed logs")
parser.add_argument("-d", "--debug", action="store_true",
help="Enable debug logging")
parser.add_argument("-f", "--follow", action="store_true",
help="Follow log output until job completes")
parser.add_argument("url", help="Base log URL (e.g. https://…/slug/log)")
args = parser.parse_args()

if args.debug:
logging.basicConfig(format="%(message)s", level=logging.DEBUG)

url = args.url
if url.endswith("/"):
url += "log"

with contextlib.suppress(KeyboardInterrupt):
for chunk in fetch_log(url, follow=args.follow):
print(chunk, end='', flush=True)


if __name__ == "__main__":
main()
Loading