Skip to content
Merged
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
1 change: 0 additions & 1 deletion Package.swift
Original file line number Diff line number Diff line change
Expand Up @@ -183,7 +183,6 @@ let package = Package(
.product(name: "GoogleAuth", package: "swift-google-auth"),
.product(name: "GoogleGax", package: "swift-google-gax"),
.product(name: "Logging", package: "swift-log"),
.product(name: "NIOCore", package: "swift-nio"),
],
path: "Tests/StorageW1R3",
exclude: ["README.md"]
Expand Down
17 changes: 7 additions & 10 deletions Tests/StorageW1R3/BenchmarkRunner.swift
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@
// limitations under the License.

import Foundation
import NIOCore
import GoogleAuth
import GoogleGax
import GoogleCloudStorage
Expand All @@ -23,7 +22,7 @@ extension StorageW1R3 {
func runWorker(
taskIndex: Int,
counters: BenchmarkCounters,
buffer: NIOCore.ByteBuffer,
buffer: ByteChunk,
storageClient: StorageClient,
controlClient: StorageControlClient,
) async {
Expand Down Expand Up @@ -55,7 +54,7 @@ extension StorageW1R3 {
let objectName = Self.randomObjectName()
let isResumable = Bool.random()
let uploadCrc32c = self.pickCrc32c()
let uploadSlice = buffer.getSlice(at: 0, length: size) ?? buffer.slice()
let uploadSlice = buffer.subdata(in: 0..<size)
let iterationId = IterationId(
task: taskIndex, taskStartInstant: taskStartInstant, iteration: iteration)

Expand Down Expand Up @@ -124,7 +123,7 @@ extension StorageW1R3 {
controlClient: StorageControlClient,
bucketName: String,
objectName: String,
buffer: NIOCore.ByteBuffer,
buffer: ByteChunk,
isResumable: Bool,
crc32cEnabled: Bool
) async -> GoogleCloudStorage.Object? {
Expand All @@ -133,7 +132,7 @@ extension StorageW1R3 {
let uploadBuilder = SampleBuilder(
iterationId: iterationId,
op: uploadOp,
targetSize: buffer.readableBytes,
targetSize: buffer.count,
object: objectName,
crc32cEnabled: crc32cEnabled
)
Expand Down Expand Up @@ -246,10 +245,9 @@ extension StorageW1R3 {
print(sample.toRow())
}

func generateRandomBuffer() -> NIOCore.ByteBuffer {
func generateRandomBuffer() -> ByteChunk {
let size = self.maxObjectSize
var buffer = ByteBufferAllocator().buffer(capacity: size)
guard size > 0 else { return buffer }
guard size > 0 else { return ByteChunk() }
// There is a lot going on here. Sometimes the benchmark is used with really large buffers,
// 256MiB and 2GiB are not uncommon. To efficiently initialized the buffer with random data
// we create an array of the desired size.
Expand All @@ -268,8 +266,7 @@ extension StorageW1R3 {
}
initializedCount = size
}
buffer.writeBytes(bytes)
return buffer
return ByteChunk(bytes)
}

private static func randomObjectName() -> String {
Expand Down
9 changes: 4 additions & 5 deletions Tests/StorageW1R3/StorageOperations.swift
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@
// limitations under the License.

import Foundation
import NIOCore
import GoogleAuth
import GoogleGax
import GoogleCloudStorage
Expand All @@ -25,7 +24,7 @@ enum StorageOperations {
controlClient: StorageControlClient,
bucketName: String,
objectName: String,
buffer: NIOCore.ByteBuffer,
buffer: ByteChunk,
isResumable: Bool,
crc32cEnabled: Bool
) async throws -> GoogleCloudStorage.Object {
Expand All @@ -37,15 +36,15 @@ enum StorageOperations {
// If resumable, chunk size is set to 32MiB; if simple, threshold handles it
if isResumable {
$0.chunkSize = 32 * 1024 * 1024
$0.resumableUploadThreshold = buffer.readableBytes
$0.resumableUploadThreshold = buffer.count
} else {
$0.resumableUploadThreshold = buffer.readableBytes + 256 * 1024
$0.resumableUploadThreshold = buffer.count + 256 * 1024
}
}

do {
return try await client.upload(
BytesSource(buffer: .init(buffer)), to: bucketName, as: objectName, options: options)
BytesSource(buffer: buffer), to: bucketName, as: objectName, options: options)
} catch let reqError as RequestError where reqError.isFailedPrecondition {
logToStderr("Precondition failed for \(objectName), fetching object details")
let getReq = GetObjectRequest().with {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,55 +17,47 @@ import NIOCore
import NIOFoundationCompat

/// A container representing a sequence of bytes backed either by `Foundation.Data`
/// or `NIOCore.ByteBuffer` without unnecessary memory copying.
public struct ByteBuffer: Sendable {
@usableFromInline
/// or an internal network buffer without unnecessary memory copying.
public struct ByteChunk: Sendable, ContiguousBytes {
internal enum Storage: Sendable {
case data(Data)
case byteBuffer(NIOCore.ByteBuffer)
}

@usableFromInline
internal let storage: Storage

// MARK: - Initializers

/// Creates a byte buffer wrapping a `Foundation.Data` instance (zero-copy).
@inlinable
/// Creates a byte chunk wrapping a `Foundation.Data` instance (zero-copy).
public init(_ data: Data) {
self.storage = .data(data)
}

/// Creates a byte buffer wrapping a `NIOCore.ByteBuffer` instance (zero-copy).
@inlinable
public init(_ buffer: NIOCore.ByteBuffer) {
/// Creates a byte chunk wrapping a `NIOCore.ByteBuffer` instance (zero-copy).
internal init(_ buffer: NIOCore.ByteBuffer) {
self.storage = .byteBuffer(buffer)
}

/// Creates an empty byte buffer instance.
@inlinable
/// Creates an empty byte chunk instance.
public init() {
self.storage = .data(Data())
self.storage = .byteBuffer(NIOCore.ByteBuffer())
}

/// Creates a byte buffer from an array of bytes.
@inlinable
/// Creates a byte chunk from an array of bytes.
public init(_ bytes: [UInt8]) {
self.storage = .data(Data(bytes))
self.storage = .byteBuffer(NIOCore.ByteBuffer(bytes: bytes))
}

/// Creates a byte buffer from a contiguous raw buffer pointer.
@inlinable
/// Creates a byte chunk from a contiguous raw buffer pointer.
public init(_ bufferPointer: UnsafeRawBufferPointer) {
self.storage = .data(Data(bufferPointer))
self.storage = .byteBuffer(NIOCore.ByteBuffer(bytes: bufferPointer))
Comment thread
coryan marked this conversation as resolved.
}
}

// MARK: - Core Properties & Accessors

extension ByteBuffer {
extension ByteChunk {
/// The total number of readable bytes stored.
@inlinable
public var count: Int {
switch storage {
case .data(let data):
Expand All @@ -75,14 +67,13 @@ extension ByteBuffer {
}
}

/// Indicates whether the buffer contains zero bytes.
/// Indicates whether the chunk contains zero bytes.
@inlinable
public var isEmpty: Bool {
count == 0
}

/// Calls a closure with a pointer to the contiguous bytes without copying.
@inlinable
public func withUnsafeBytes<R>(_ body: (UnsafeRawBufferPointer) throws -> R) rethrows -> R {
switch storage {
case .data(let data):
Expand All @@ -92,11 +83,22 @@ extension ByteBuffer {
}
}

/// Executes a closure on the sequence's contiguous storage.
@inlinable
public func withContiguousStorageIfAvailable<R>(
_ body: (UnsafeBufferPointer<UInt8>) throws -> R
) rethrows -> R? {
try withUnsafeBytes { rawBuffer in
try rawBuffer.withMemoryRebound(to: UInt8.self) { buffer in
try body(buffer)
}
}
}

/// The underlying contents as a `Foundation.Data` instance.
///
/// - Returns: The original `Data` with zero copies if backed by `Data`,
/// or copies the bytes into a new `Data` instance if backed by `NIOCore.ByteBuffer`.
@inlinable
/// or copies the bytes into a new `Data` instance if backed by an internal network buffer.
public var data: Data {
switch storage {
case .data(let data):
Expand All @@ -110,8 +112,7 @@ extension ByteBuffer {
///
/// - Returns: The original `NIOCore.ByteBuffer` with zero copies if backed by `NIOCore.ByteBuffer`,
/// or copies the bytes into a new `NIOCore.ByteBuffer` instance if backed by `Data`.
@inlinable
public var byteBuffer: NIOCore.ByteBuffer {
internal var byteBuffer: NIOCore.ByteBuffer {
switch storage {
case .byteBuffer(let buffer):
return buffer
Expand All @@ -130,27 +131,27 @@ extension ByteBuffer {
withUnsafeBytes { Array($0) }
}

/// Returns a zero-copy sub-buffer within the specified byte range.
public func subdata(in range: Range<Int>) -> ByteBuffer {
/// Returns a zero-copy sub-chunk within the specified byte range.
public func subdata(in range: Range<Int>) -> ByteChunk {
switch storage {
case .data(let data):
let start = data.startIndex.advanced(by: range.lowerBound)
let end = data.startIndex.advanced(by: range.upperBound)
return ByteBuffer(data[start..<end])
return ByteChunk(data[start..<end])
case .byteBuffer(let nioBuffer):
var copy = nioBuffer
copy.moveReaderIndex(to: nioBuffer.readerIndex + range.lowerBound)
if let slice = copy.readSlice(length: range.count) {
return ByteBuffer(slice)
return ByteChunk(slice)
}
return ByteBuffer()
return ByteChunk()
}
}
}

// MARK: - RandomAccessCollection Conformance

extension ByteBuffer: RandomAccessCollection {
extension ByteChunk: RandomAccessCollection {
public typealias Element = UInt8
public typealias Index = Int

Expand All @@ -160,7 +161,6 @@ extension ByteBuffer: RandomAccessCollection {
@inlinable
public var endIndex: Int { count }

@inlinable
public subscript(position: Int) -> UInt8 {
precondition(position >= 0 && position < count, "Index \(position) out of bounds 0..<\(count)")
switch storage {
Expand All @@ -174,8 +174,8 @@ extension ByteBuffer: RandomAccessCollection {

// MARK: - Equatable & Hashable

extension ByteBuffer: Equatable {
public static func == (lhs: ByteBuffer, rhs: ByteBuffer) -> Bool {
extension ByteChunk: Equatable {
public static func == (lhs: ByteChunk, rhs: ByteChunk) -> Bool {
guard lhs.count == rhs.count else { return false }
if lhs.isEmpty { return true }
return lhs.withUnsafeBytes { lhsBytes in
Expand All @@ -189,21 +189,22 @@ extension ByteBuffer: Equatable {
}
}

extension ByteBuffer: Hashable {
extension ByteChunk: Hashable {
@inlinable
public func hash(into hasher: inout Hasher) {
withUnsafeBytes { hasher.combine(bytes: $0) }
}
}

// MARK: - Literal & Description Conformances

extension ByteBuffer: ExpressibleByArrayLiteral {
extension ByteChunk: ExpressibleByArrayLiteral {
public init(arrayLiteral elements: UInt8...) {
self.init(Data(elements))
self.init(elements)
}
}

extension ByteBuffer: CustomStringConvertible, CustomDebugStringConvertible {
extension ByteChunk: CustomStringConvertible, CustomDebugStringConvertible {
public var description: String {
"\(count) bytes"
}
Expand All @@ -214,6 +215,6 @@ extension ByteBuffer: CustomStringConvertible, CustomDebugStringConvertible {
case .data: backing = "Data"
case .byteBuffer: backing = "NIOCore.ByteBuffer"
}
return "ByteBuffer(\(count) bytes, backing: \(backing))"
return "ByteChunk(\(count) bytes, backing: \(backing))"
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -13,25 +13,28 @@
// limitations under the License.

import Foundation
import NIOCore

/// An upload source that wraps in-memory bytes or buffers.
public struct BytesSource: SeekableUploadSource {
public let buffer: ByteBuffer
public let buffer: ByteChunk
public var totalSize: UInt64? {
return UInt64(buffer.count)
}
private var offset: UInt64 = 0

public init(buffer: ByteBuffer) {
public init(buffer: ByteChunk) {
self.buffer = buffer
}

public init(_ chunk: ByteChunk) {
self.buffer = chunk
}

public init(data: Data) {
self.buffer = ByteBuffer(data)
self.buffer = ByteChunk(data)
}

public mutating func read(maxBytes: Int) async throws -> ByteBuffer? {
public mutating func read(maxBytes: Int) async throws -> ByteChunk? {
guard maxBytes > 0, offset < UInt64(buffer.count) else { return nil }
let end = min(offset + UInt64(maxBytes), UInt64(buffer.count))
let chunk = buffer.subdata(in: Int(offset)..<Int(end))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,8 +35,8 @@ extension ChecksumCalculator {
data.withUnsafeBytes { update($0) }
}

/// Convenience helper for ByteBuffer chunks.
mutating func update(_ buffer: ByteBuffer) {
/// Convenience helper for ByteChunk chunks.
mutating func update(_ buffer: ByteChunk) {
buffer.withUnsafeBytes { update($0) }
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
import Foundation

struct ChunkInfo: Sendable {
let data: ByteBuffer
let data: ByteChunk
let isLast: Bool
let checksum: String?
}
Expand All @@ -24,7 +24,7 @@ struct ChecksummedSource<S: UploadSource> {
var source: S
let options: ChecksumOptions
private var calculators: [any ChecksumCalculator] = []
private var nextChunk: ByteBuffer? = nil
private var nextChunk: ByteChunk? = nil
private var isInitialized = false
private var isFinished = false
/// The high-water mark of sequentially processed bytes in `calculators`.
Expand Down Expand Up @@ -78,13 +78,13 @@ struct ChecksummedSource<S: UploadSource> {
/// To support seeking backward and retrying chunk uploads without corrupting checksums,
/// this method skips any prefix of `data` that falls below `bytesHashed` (the high-water mark
/// of bytes already fed into `calculators`). Only bytes beyond `bytesHashed` are accumulated.
private mutating func updateChecksums(data: ByteBuffer, startOffset: UInt64) {
private mutating func updateChecksums(data: ByteChunk, startOffset: UInt64) {
guard !calculators.isEmpty else { return }

let endOffset = startOffset + UInt64(data.count)
guard endOffset > bytesHashed else { return }

let unhashedData: ByteBuffer
let unhashedData: ByteChunk
if startOffset >= bytesHashed {
unhashedData = data
} else {
Expand Down
Loading
Loading