diff --git a/Package.swift b/Package.swift index 330890a..8f1adae 100644 --- a/Package.swift +++ b/Package.swift @@ -24,6 +24,9 @@ let strictConcurrencySettings: [SwiftSetting] = { // -warnings-as-errors here is a workaround so that IDE-based development can // get tripped up on -require-explicit-sendable. initialSettings.append(.unsafeFlags(["-Xfrontend", "-require-explicit-sendable", "-warnings-as-errors"])) + initialSettings.append(.enableExperimentalFeature("LifetimeDependence")) + initialSettings.append(.enableExperimentalFeature("Lifetimes")) + initialSettings.append(.enableUpcomingFeature("LifetimeDependence")) } return initialSettings @@ -45,6 +48,7 @@ let package = Package( .package(url: "https://github.com/apple/swift-algorithms.git", from: "1.0.0"), .package(url: "https://github.com/apple/swift-distributed-tracing.git", from: "1.3.0"), .package(url: "https://github.com/apple/swift-configuration.git", from: "1.0.0"), + .package(path: "../swift-http-client-server-apis"), ], targets: [ .target( @@ -74,6 +78,9 @@ let package = Package( // Observability support .product(name: "Logging", package: "swift-log"), .product(name: "Tracing", package: "swift-distributed-tracing"), + + // HTTP APIs + .product(name: "HTTPAPIs", package: "swift-http-client-server-apis"), ], swiftSettings: strictConcurrencySettings ), diff --git a/Sources/AsyncHTTPClient/HTTP APIs/AHC+HTTP.swift b/Sources/AsyncHTTPClient/HTTP APIs/AHC+HTTP.swift new file mode 100644 index 0000000..bbda947 --- /dev/null +++ b/Sources/AsyncHTTPClient/HTTP APIs/AHC+HTTP.swift @@ -0,0 +1,226 @@ +//===----------------------------------------------------------------------===// +// +// This source file is part of the AsyncHTTPClient open source project +// +// Copyright (c) 2025 Apple Inc. and the AsyncHTTPClient project authors +// Licensed under Apache License v2.0 +// +// See LICENSE.txt for license information +// See CONTRIBUTORS.txt for the list of AsyncHTTPClient project authors +// +// SPDX-License-Identifier: Apache-2.0 +// +//===----------------------------------------------------------------------===// + +#if compiler(>=6.2) +import HTTPAPIs +import HTTPTypes +import NIOHTTP1 +import Foundation +import NIOCore +import Synchronization +import BasicContainers + +@available(macOS 26.2, iOS 26.2, watchOS 26.2, tvOS 26.2, *) +extension AsyncHTTPClient.HTTPClient: HTTPAPIs.HTTPClient { + public typealias RequestConcludingWriter = RequestWriter + public typealias ResponseConcludingReader = ResponseReader + + public struct RequestWriter: ConcludingAsyncWriter, ~Copyable, SendableMetatype { + public typealias Underlying = RequestBodyWriter + public typealias FinalElement = HTTPFields? + + final class Storage: Sendable { + struct StateMachine: ~Copyable { + + enum State: ~Copyable { + case buffer(CheckedContinuation?) + case demand(CheckedContinuation) + case done(HTTPFields?) + } + + var state: State = .buffer(nil) + + init() { + + } + + } + + let stateMutex = Mutex(StateMachine()) + + nonisolated(nonsending) func write( + _ body: nonisolated(nonsending) (inout OutputSpan) async throws(Failure) -> Result) async throws(AsyncStreaming.EitherError + ) -> Result where Failure : Error { + + fatalError() +// do { +// +// self.stateMutex.withLock { state in +// +// } +// +// let updated = try await self.buffer.edit { (outputSpan) async throws(Failure) -> Result in +// try await body(&outputSpan) +// } +// +// +// +// return updated +// } catch { +// throw .first(error) +// } + } + + func next() -> ByteBuffer { + fatalError() + } + } + + let storage: Storage + + init() { + self.storage = .init() + } + + public consuming func produceAndConclude( + body: nonisolated(nonsending) (consuming sending HTTPClient.RequestBodyWriter) async throws -> (Return, HTTPFields?) + ) async throws -> Return { + let bodyWriter = RequestBodyWriter(storage: self.storage) + do { + let (ret, fields) = try await body(bodyWriter) + return ret + } catch { + throw error + } + } + } + + public struct RequestBodyWriter: AsyncWriter, ~Copyable { + public typealias WriteElement = UInt8 + public typealias WriteFailure = any Error + + let storage: RequestWriter.Storage + + public mutating func write( + _ body: nonisolated(nonsending) (inout OutputSpan) async throws(Failure) -> Result) async throws(AsyncStreaming.EitherError + ) -> Result where Failure : Error { + try await self.storage.write(body) + } + } + + public struct ResponseReader: ConcludingAsyncReader { + public typealias Underlying = ResponseBodyReader + + let underlying: HTTPClientResponse.Body + + public typealias FinalElement = HTTPFields? + + init(underlying: HTTPClientResponse.Body) { + self.underlying = underlying + } + + public consuming func consumeAndConclude( + body: nonisolated(nonsending) (consuming sending HTTPClient.ResponseBodyReader) async throws(Failure) -> Return + ) async throws(Failure) -> (Return, HTTPFields?) where Failure : Error { + let iterator = self.underlying.makeAsyncIterator() + let reader = ResponseBodyReader(underlying: iterator) + let returnValue = try await body(reader) + return (returnValue, nil) + } + + } + + public struct ResponseBodyReader: AsyncReader, ~Copyable { + public typealias ReadElement = UInt8 + public typealias ReadFailure = any Error + + var underlying: HTTPClientResponse.Body.AsyncIterator + + public mutating func read( + maximumCount: Int?, + body: nonisolated(nonsending) (consuming Span) async throws(Failure) -> Return + ) async throws(AsyncStreaming.EitherError) -> Return where Failure : Error { + + do { + let buffer = try await self.underlying.next(isolation: #isolation) + if let buffer { + var array = RigidArray() + array.reserveCapacity(buffer.readableBytes) + buffer.withUnsafeReadableBytes { rawBufferPtr in + let usbptr = rawBufferPtr.assumingMemoryBound(to: UInt8.self) + array.append(copying: usbptr) + } + return try await body(array.span) + } else { + let array = InlineArray<0, UInt8> { _ in } + return try await body(array.span) + } + } catch let error as Failure { + throw .second(error) + } catch { + throw .first(error) + } + } + } + + public func perform( + request: HTTPRequest, + body: consuming HTTPClientRequestBody?, + configuration: HTTPClientConfiguration, + eventHandler: borrowing some HTTPClientEventHandler & ~Copyable & ~Escapable, + responseHandler: nonisolated(nonsending) (HTTPResponse, consuming ResponseReader) async throws -> Return + ) async throws -> Return { + guard let url = request.url else { + fatalError() + } + + var ahcRequest = HTTPClientRequest(url: url.absoluteString) + ahcRequest.method = .init(rawValue: request.method.rawValue) + if !request.headerFields.isEmpty { + let sequence = request.headerFields.lazy.map({ ($0.name.rawName, $0.value) }) + ahcRequest.headers.add(contentsOf: sequence) + } + + let result = try await withThrowingTaskGroup { taskGroup in + switch body { + case .none: + break + + case .restartable(let handler): + taskGroup.addTask { + let writer = RequestWriter() + try await handler(writer) + } + case .some(.seekable(_)): + fatalError() + } + + let ahcResponse = try await self.execute(ahcRequest, timeout: .seconds(30)) + + var responseFields = HTTPFields() + for (name, value) in ahcResponse.headers { + if let name = HTTPField.Name(name) { + responseFields[name] = value + } + } + + let response = HTTPResponse( + status: .init(code: Int(ahcResponse.status.code)), + headerFields: responseFields + ) + + return try await responseHandler(response, .init(underlying: ahcResponse.body)) + + } + + return result + } +} + +//private struct ClosureAsyncSequence: AsyncSequence { +// var body: +// +//} + +#endif diff --git a/Tests/AsyncHTTPClientTests/HTTP APIs/AHC+HTTP_Tests.swift b/Tests/AsyncHTTPClientTests/HTTP APIs/AHC+HTTP_Tests.swift new file mode 100644 index 0000000..9c5e681 --- /dev/null +++ b/Tests/AsyncHTTPClientTests/HTTP APIs/AHC+HTTP_Tests.swift @@ -0,0 +1,82 @@ +// +// AHC+HTTP_Tests.swift +// async-http-client +// +// Created by Fabian Fett on 02.12.25. +// + +#if compiler(>=6.2) +import NIOCore +import HTTPTypes +import HTTPAPIs +import Testing +import AsyncHTTPClient + +#if canImport(Darwin) +public import Security +#endif + +@Suite +struct AbstractHTTPClientTest { + + @available(macOS 26.2, iOS 26.2, watchOS 26.2, tvOS 26.2, visionOS 26.2, *) + struct EventHandler: HTTPClientEventHandler { + func handleRedirection(response: HTTPTypes.HTTPResponse, newRequest: HTTPTypes.HTTPRequest) async throws -> HTTPAPIs.HTTPClientRedirectionAction { + .deliverRedirectionResponse + } + + #if canImport(Darwin) + func handleServerTrust(_ trust: SecTrust) async throws -> HTTPAPIs.HTTPClientTrustResult { + fatalError() + } + #endif + } + + @available(macOS 26.2, iOS 26.2, watchOS 26.2, tvOS 26.2, visionOS 26.2, *) + @Test func testExample() async throws { + + let bin = HTTPBin(.http1_1(ssl: false)) + defer { try! bin.shutdown() } + + let client = HTTPClient() + defer { try! client.shutdown().wait() } + + let request = HTTPRequest(method: .get, scheme: "http", authority: "127.0.0.1:\(bin.port)", path: "/stats") +// let body = HTTPClientRequestBody.restartable { (writer: consuming AsyncHTTPClient.HTTPClient.RequestWriter) in +// try await writer.produceAndConclude { writer in +// +// var mwriter = writer +// +// try await mwriter.write { outputSpan in +// outputSpan.append(repeating: UInt8(ascii: "X"), count: 10) +// } +// +// return ((), nil) +// } +// } + + let handler = EventHandler() + + try await client.perform( + request: request, + body: nil, + configuration: .init(), + eventHandler: handler + ) { response, responseReader in + print("status: \(response.status)") + for header in response.headerFields { + print("\(header.name): \(header.value)") + } + + let trailers = try await responseReader.collect(upTo: 1024) { span in + span.withUnsafeBufferPointer { buffer in + print(String(decoding: buffer, as: Unicode.UTF8.self)) + } + } + } + + } + +} + +#endif