Everett
Loading...
Searching...
No Matches
object_stream.h
Go to the documentation of this file.
1
13#pragma once
14
16
17#include <optional>
18
19namespace everett {
20 // Append an object whose final extent is not yet known. Each append consumes
21 // complete physical bytes, including for bit policies. The caller keeps an
22 // unfinished bit tail private until it can supply the final canonical byte.
23 // Body chunks need live only during append. An optional fixed-size prefix is
24 // initially zero and supplied at finish, for directories known only afterward.
25 // Its checksum combines with the streamed suffix without rereading payloads.
26 //
27 // Construction reserves an exclusive private file under caller-reserved
28 // identities. The root and filesystem premises are object_writer's. No
29 // partial-output persistence receipt or restart operation is provided.
30 // Destruction closes handles and preserves surviving names after failure or
31 // abandonment. A borrowed Ops must outlive this nonmovable stream.
32 template <class P, class Ops = posix_object_ops> struct object_stream {
33 static_assert(Ops::supported, "Everett object streams require supported filesystem operations");
34 using policy_type = P;
35
36 object_stream(std::filesystem::path root, object_id id,
37 object_attempt_id attempt, file_kind kind, std::size_t prefix_bytes = 0)
38 : root_(std::move(root)), id_(std::move(id)), attempt_(std::move(attempt)),
39 kind_(checked(kind)), prefix_bytes_(checked_prefix(prefix_bytes)),
40 owned_ops_(std::in_place), ops_(*owned_ops_),
41 run_(root_, id_, attempt_, kind_, ops_) { open(); }
42 object_stream(std::filesystem::path root, object_id id,
43 object_attempt_id attempt, file_kind kind, Ops & ops)
44 : object_stream(std::move(root), std::move(id), std::move(attempt), kind, 0, ops) {}
45 object_stream(std::filesystem::path root, object_id id,
46 object_attempt_id attempt, file_kind kind, std::size_t prefix_bytes, Ops & ops)
47 : root_(std::move(root)), id_(std::move(id)), attempt_(std::move(attempt)),
48 kind_(checked(kind)), prefix_bytes_(checked_prefix(prefix_bytes)),
49 ops_(ops), run_(root_, id_, attempt_, kind_, ops_) { open(); }
50 object_stream(object_stream const &) = delete;
54
55 std::uint64_t body_bytes() const noexcept { return body_bytes_; }
56 // Before finish this describes the zero-filled prefix and appended bytes;
57 // afterward it describes the final prefix and body. Failure makes no claim
58 // about which writes reached the underlying file.
59 std::uint32_t body_crc32c() const noexcept { return crc_; }
60 bool failed() const noexcept { return failed_; }
61 bool finished() const noexcept { return finished_; }
62 object_write_paths const & paths() const & noexcept { return run_.paths; }
63 object_write_paths const & paths() const && = delete;
64
65 // Invalid size is rejected before writing. Once a write fails, acknowledged
66 // byte/CRC progress cannot certify the uncertain file and the stream is
67 // poisoned. A later successful call must never rehabilitate that attempt.
68 void append(std::span<std::byte const> bytes) {
70 constexpr auto maximum = std::uint64_t(std::numeric_limits<std::int64_t>::max()) -
72 if (bytes.size() > maximum - body_bytes_)
73 throw std::length_error("Everett object exceeds supported file offsets");
74 try {
75 while (!bytes.empty()) {
76 auto part = bytes.first(std::min(bytes.size(), std::size_t{1} << 20));
77 run_.write_all(part);
78 crc_ = crc32c<typename P::architecture>(part, crc_);
79 body_bytes_ += part.size();
80 last_ = part.back();
81 bytes = bytes.subspan(part.size());
82 }
83 } catch (...) {
84 failed_ = true;
85 throw;
86 }
87 }
88
89 // Metadata-only rejection leaves the stream available for corrected
90 // metadata or more body bytes. Once sealing starts, any failure poisons it.
92 return finish(header, {});
93 }
94 object_seal_receipt finish(file_header<P> const & header, std::span<std::byte const> prefix) {
97 if (prefix.size() != prefix_bytes_ || header.kind != kind_ ||
98 file_detail::body_bytes<P>(header.extent) != body_bytes_)
99 throw std::invalid_argument("Everett streamed body metadata mismatch");
100 auto last = body_bytes_ == prefix_bytes_ && !prefix.empty() ? prefix.back() : last_;
101 if constexpr (P::unit == profile_unit::bit)
102 if ((header.extent & 7) &&
103 (std::to_integer<unsigned>(last) & ((1u << (8 - (header.extent & 7))) - 1)))
104 throw std::invalid_argument("noncanonical bit-profile tail padding");
105 auto crc = crc_ ^ crc32c_combine(prefix_crc_ ^ crc32c<typename P::architecture>(prefix), 0, body_bytes_ - prefix_bytes_);
106 auto encoded = encode_file_header(header, crc);
107 object_seal_receipt receipt{id_, attempt_, run_.paths.final,
108 body_bytes_ + file_detail::header_bytes, crc, ops_.barrier()};
109 try {
110 run_.write_at(prefix, file_detail::header_bytes, "write object prefix");
111 run_.write_header(encoded);
113 run_.finish();
114 crc_ = crc;
115 finished_ = true;
116 return receipt;
117 } catch (...) {
118 failed_ = true;
119 throw;
120 }
121 }
122
123 private:
124 static std::size_t checked_prefix(std::size_t bytes) {
125 if (bytes > std::uint64_t(std::numeric_limits<std::int64_t>::max()) - file_detail::header_bytes)
126 throw std::length_error("Everett object prefix exceeds supported file offsets");
127 return bytes;
128 }
131 throw std::invalid_argument("unsupported Everett object kind");
132 return kind;
133 }
134 void open() {
135 run_.open();
136 std::array<std::byte, file_detail::header_bytes> placeholder{};
137 run_.write_all(placeholder);
138 std::array<std::byte, 4096> zero{};
139 auto left = prefix_bytes_;
140 while (left) {
141 auto count = std::min(left, zero.size());
142 append(std::span(zero).first(count));
143 left -= count;
144 }
146 }
147 void require_active() const {
148 if (failed_ || finished_) throw std::logic_error("Everett object stream is inactive");
149 }
150 std::filesystem::path root_;
154 std::size_t prefix_bytes_ = 0;
155 std::optional<Ops> owned_ops_;
156 Ops & ops_;
158 std::uint64_t body_bytes_ = 0;
159 std::uint32_t crc_ = 0;
160 std::uint32_t prefix_crc_ = 0;
161 std::byte last_{};
162 bool failed_ = false;
163 bool finished_ = false;
164 };
165}
void validate_metadata(file_header< P > const &header)
Definition file.h:93
constexpr std::size_t header_bytes
Definition file.h:47
Definition active_engine.h:18
std::uint32_t crc32c_combine(std::uint32_t first, std::uint32_t second, std::uint64_t second_bytes) noexcept
std::array< std::byte, file_detail::header_bytes > encode_file_header(file_header< P > const &header, std::uint32_t body_crc)
Definition file.h:194
file_kind
Definition object_path.h:25
Streams and seals immutable object files without catalog publication.
Definition file.h:34
file_kind kind
Definition file.h:36
std::uint64_t extent
Definition file.h:37
Definition object_writer.h:39
Definition object_path.h:38
Definition object_writer.h:88
Definition object_stream.h:32
void append(std::span< std::byte const > bytes)
Definition object_stream.h:68
std::uint64_t body_bytes_
Definition object_stream.h:158
object_id id_
Definition object_stream.h:151
std::optional< Ops > owned_ops_
Definition object_stream.h:155
std::byte last_
Definition object_stream.h:161
Ops & ops_
Definition object_stream.h:156
std::uint32_t crc_
Definition object_stream.h:159
bool failed_
Definition object_stream.h:162
std::uint32_t prefix_crc_
Definition object_stream.h:160
object_write_paths const & paths() const &noexcept
Definition object_stream.h:62
object_writer< P, Ops >::operation run_
Definition object_stream.h:157
void open()
Definition object_stream.h:134
object_stream & operator=(object_stream const &)=delete
std::uint64_t body_bytes() const noexcept
Definition object_stream.h:55
object_stream & operator=(object_stream &&)=delete
object_stream(std::filesystem::path root, object_id id, object_attempt_id attempt, file_kind kind, std::size_t prefix_bytes=0)
Definition object_stream.h:36
bool failed() const noexcept
Definition object_stream.h:60
std::uint32_t body_crc32c() const noexcept
Definition object_stream.h:59
std::filesystem::path root_
Definition object_stream.h:150
bool finished_
Definition object_stream.h:163
object_attempt_id attempt_
Definition object_stream.h:152
object_stream(object_stream &&)=delete
void require_active() const
Definition object_stream.h:147
object_seal_receipt finish(file_header< P > const &header, std::span< std::byte const > prefix)
Definition object_stream.h:94
object_stream(std::filesystem::path root, object_id id, object_attempt_id attempt, file_kind kind, std::size_t prefix_bytes, Ops &ops)
Definition object_stream.h:45
file_kind kind_
Definition object_stream.h:153
object_stream(object_stream const &)=delete
bool finished() const noexcept
Definition object_stream.h:61
P policy_type
Definition object_stream.h:34
static std::size_t checked_prefix(std::size_t bytes)
Definition object_stream.h:124
static file_kind checked(file_kind kind)
Definition object_stream.h:129
object_stream(std::filesystem::path root, object_id id, object_attempt_id attempt, file_kind kind, Ops &ops)
Definition object_stream.h:42
object_write_paths const & paths() const &&=delete
object_seal_receipt finish(file_header< P > const &header)
Definition object_stream.h:91
std::size_t prefix_bytes_
Definition object_stream.h:154
Definition object_writer.h:53
Definition object_writer.h:235