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
13 changes: 13 additions & 0 deletions pkgs/cronet_http/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,16 @@
## 1.10.0-wip

* Add support for streaming request bodies on `StreamedRequest` via Cronet's
`UploadDataProvider` API. In-memory `Request` bodies continue to use the
existing byte-buffer upload path.
* Add `UploadDataProviderProxy` Kotlin bridge so Dart can implement Cronet's
`UploadDataProvider` through JNI (jnigen cannot subclass abstract Java
classes).
* Regenerate JNI bindings for `UploadDataProvider`, `UploadDataSink`, and
`UploadDataProviderProxy`.
* Add `package:async` dependency for `StreamQueue` when streaming uploads.
* Enable streamed request body conformance tests (`canStreamRequestBody: true`).

## 1.9.0

* Add `CronetEngine.startNetLogToFile` and `CronetEngine.stopNetLog`.
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
// Copyright (c) 2023, the Dart project authors. Please see the AUTHORS file
// for details. All rights reserved. Use of this source code is governed by a
// BSD-style license that can be found in the LICENSE file.

// Cronet uploads request bodies by subclassing the abstract class
// `UploadDataProvider`. Cronet calls `getLength()`, then repeatedly calls
// `read()` to pull bytes into a `ByteBuffer`, and may call `rewind()` on
// redirects before calling `close()` when the upload finishes.
//
// `package:jnigen` does not support subclassing abstract Java classes from Dart
// (see https://github.com/dart-lang/jnigen/issues/348).
//
// This file provides an interface `UploadDataProviderInterface`, which can be
// implemented in Dart, and a wrapper class `UploadDataProviderProxy`, which
// can be passed to Cronet's `UrlRequest.Builder.setUploadDataProvider`.

package io.flutter.plugins.cronet_http

import androidx.annotation.Keep
import org.chromium.net.UploadDataProvider
import org.chromium.net.UploadDataSink
import java.nio.ByteBuffer

// Due to a bug (https://github.com/dart-lang/native/issues/2421) where JNIgen
// does not synchronize nullabilities across the class hierarchy and the fact
// that UploadDataProvider is a Java class with no nullability annotations,
// generating both `UploadDataProviderProxy` and `UploadDataProvider` together
// with different nullabilities causes the super method to have a looser type
// for parameters which is a Dart compilation error.
// That is why `read` and `rewind` parameters on the interface are defined as
// nullable to match `UploadDataProvider` while in reality Cronet always
// passes non-null values.

@Keep
class UploadDataProviderProxy(
private val callback: UploadDataProviderInterface
) : UploadDataProvider() {

@Keep
interface UploadDataProviderInterface {
fun getLength(): Long
fun read(uploadDataSink: UploadDataSink?, byteBuffer: ByteBuffer?)
fun rewind(uploadDataSink: UploadDataSink?)
fun close()
}

override fun getLength(): Long = callback.getLength()

override fun read(uploadDataSink: UploadDataSink, byteBuffer: ByteBuffer) =
callback.read(uploadDataSink, byteBuffer)

override fun rewind(uploadDataSink: UploadDataSink) =
callback.rewind(uploadDataSink)

override fun close() = callback.close()
}
6 changes: 3 additions & 3 deletions pkgs/cronet_http/example/integration_test/client_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ Future<void> testConformance() async {
try {
testAll(
CronetClient.defaultCronetEngine,
canStreamRequestBody: false,
canStreamRequestBody: true,
canReceiveSetCookieHeaders: true,
canSendCookieHeaders: true,
supportsAbort: true,
Expand All @@ -34,7 +34,7 @@ Future<void> testConformance() async {
try {
testAll(
CronetClient.defaultCronetEngine,
canStreamRequestBody: false,
canStreamRequestBody: true,
canReceiveSetCookieHeaders: true,
canSendCookieHeaders: true,
supportsAbort: true,
Expand All @@ -52,7 +52,7 @@ Future<void> testConformance() async {
cacheMode: CacheMode.disabled, userAgent: 'Test Agent (Future)');
return CronetClient.fromCronetEngine(engine);
},
canStreamRequestBody: false,
canStreamRequestBody: true,
canReceiveSetCookieHeaders: true,
canSendCookieHeaders: true,
supportsAbort: true,
Expand Down
3 changes: 3 additions & 0 deletions pkgs/cronet_http/jnigen.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ output:

classes:
- 'io.flutter.plugins.cronet_http.UrlRequestCallbackProxy'
- 'io.flutter.plugins.cronet_http.UploadDataProviderProxy'
- 'java.io.IOException'
- 'java.lang.Exception'
- 'java.lang.Throwable'
Expand All @@ -24,3 +25,5 @@ classes:
- 'org.chromium.net.UploadDataProviders'
- 'org.chromium.net.UrlRequest'
- 'org.chromium.net.UrlResponseInfo'
- 'org.chromium.net.UploadDataSink'
- 'org.chromium.net.UploadDataProvider'
175 changes: 154 additions & 21 deletions pkgs/cronet_http/lib/src/cronet_client.dart
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,10 @@
// BSD-style license that can be found in the LICENSE file.

import 'dart:async';
import 'dart:math';
import 'dart:typed_data';

import 'package:async/async.dart';
import 'package:http/http.dart';
import 'package:http_profile/http_profile.dart';
import 'package:jni/jni.dart';
Expand Down Expand Up @@ -650,6 +653,113 @@ class CronetClient extends BaseClient {
requestMethod: request.method,
requestUri: request.url.toString());

/// Returns true if [stream] includes at least one list with an element.
///
/// Since [_hasData] consumes [stream], returns a new stream containing the
/// equivalent data.
static Future<(bool, Stream<List<int>>)> _hasData(
Stream<List<int>> stream,
) async {
final queue = StreamQueue(stream);
while (await queue.hasNext && (await queue.peek).isEmpty) {
await queue.next;
}

return (await queue.hasNext, queue.rest);
}

/// Streams [stream] to Cronet on demand.
(jb.UploadDataProvider, Future<void> Function()) _streamingUploadProvider(
Stream<List<int>> stream,
int? contentLength,
HttpClientRequestProfile? profile,
) {
// Cronet's UploadDataProvider.read() rejects a non-final zero-byte read, so
// strip empty chunks (a StreamedRequest may emit them anywhere) to keep
// every read non-empty. The final zero-byte read is handled
// via onReadSucceeded(true).
final queue = StreamQueue<List<int>>(stream.where((c) => c.isNotEmpty));
Uint8List? current;
var offset = 0;
var bytesSent = 0;
var disposed = false;
Future<void> dispose() async {
if (disposed) return;
disposed = true;
await queue.cancel(immediate: true);
}

// true if `current` has unconsumed bytes; false at end of stream.
Future<bool> ensureChunk() async {
if (current != null && offset < current!.length) return true;
if (!await queue.hasNext) {
current = null;
return false;
}
final chunk = await queue.next;
current = chunk is Uint8List ? chunk : Uint8List.fromList(chunk);
offset = 0;
return true;
}

final impl =
jb.UploadDataProviderProxy$UploadDataProviderInterface.implement(
jb.$UploadDataProviderProxy$UploadDataProviderInterface(
getLength: () => contentLength ?? -1,
read$async: true,
read: (uploadDataSink, byteBuffer) async {
final sink = uploadDataSink!;
try {
if (!await ensureChunk()) {
if (contentLength == null) {
sink.onReadSucceeded(true);
} else {
sink.onReadError(jb.IOException.new1(
'Body ended before contentLength'.toJString()));
}
return;
}

final dst = byteBuffer!;
final pos = dst.position;
final available = current!.length - offset;

if (contentLength != null &&
available > contentLength - bytesSent) {
sink.onReadError(jb.IOException.new1(
'Body exceeded contentLength'.toJString(),
));
return;
}

final n = min(dst.remaining, available);
dst.asUint8List().setRange(pos, pos + n, current!, offset);
dst.position = pos + n;

profile?.requestData.bodySink.add(
Uint8List.sublistView(current!, offset, offset + n),
);

offset += n;
bytesSent += n;
sink.onReadSucceeded(false);
} catch (e) {
sink.onReadError(jb.IOException.new1('$e'.toJString()));
}
},
rewind: (uploadDataSink) {
// One-shot stream: cannot replay.
uploadDataSink!.onRewindError(jb.IOException.new1(
'Streamed request bodies cannot be rewound'.toJString()));
},
close: () {
unawaited(dispose());
},
),
);
return (jb.UploadDataProviderProxy(impl), dispose);
}

/// Sends an HTTP request and asynchronously returns the response.
@override
Future<CronetStreamedResponse> send(BaseRequest request) async {
Expand Down Expand Up @@ -684,10 +794,22 @@ class CronetClient extends BaseClient {
}

final stream = request.finalize();
final body = await stream.toBytes();
profile?.requestData.bodySink.add(body);

final inMemoryBody = request is Request ? await stream.toBytes() : null;
if (inMemoryBody != null) {
profile?.requestData.bodySink.add(inMemoryBody);
}

final bool hasBody;
Stream<List<int>>? bodyStream;
if (inMemoryBody != null) {
hasBody = inMemoryBody.isNotEmpty;
} else {
(hasBody, bodyStream) = await _hasData(stream);
}

final responseCompleter = Completer<CronetStreamedResponse>();
Future<void> Function()? disposeUpload;

return await using((arena) async {
final jUrl = request.url.toString().toJString()..releasedBy(arena);
Expand All @@ -702,7 +824,7 @@ class CronetClient extends BaseClient {
..setHttpMethod(jMethod);

var headers = request.headers;
if (body.isNotEmpty &&
if (hasBody &&
!headers.keys.any((h) => h.toLowerCase() == 'content-type')) {
// Cronet requires that requests containing upload data set a
// 'Content-Type' header.
Expand All @@ -711,33 +833,44 @@ class CronetClient extends BaseClient {
headers.forEach((k, v) => builder.addHeader(
k.toJString()..releasedBy(arena), v.toJString()..releasedBy(arena)));

if (body.isNotEmpty) {
final JByteBuffer data;
try {
data = body.toJByteBuffer()..releasedBy(arena);
} on JThrowable catch (e) {
// There are no unit tests for this code. You can verify this behavior
// manually by incrementally increasing the amount of body data in
// `CronetClient.post` until you get this exception.
if (e.message.contains('java.lang.OutOfMemoryError:')) {
throw ClientException(
'Not enough memory for request body: ${e.message}',
request.url);
if (hasBody) {
if (inMemoryBody != null) {
final JByteBuffer data;
try {
data = inMemoryBody.toJByteBuffer()..releasedBy(arena);
} on JThrowable catch (e) {
// There are no unit tests for this code. You can verify this
// behavior manually by incrementally increasing the amount of body
// data in `CronetClient.post` until you get this exception.
if (e.message.contains('java.lang.OutOfMemoryError:')) {
throw ClientException(
'Not enough memory for request body: ${e.message}',
request.url);
}
rethrow;
}
rethrow;
builder.setUploadDataProvider(
jb.UploadDataProviders.create$2(data), _executor);
} else {
final (provider, dispose) = _streamingUploadProvider(
bodyStream!, request.contentLength, profile);
disposeUpload = dispose;
builder.setUploadDataProvider(provider, _executor);
}

builder.setUploadDataProvider(
jb.UploadDataProviders.create$2(data), _executor);
}

// Not releasing `cronetRequest` as it's used in `whenComplete` callback.
final cronetRequest = builder.build()!;
if (request case Abortable(:final abortTrigger?)) {
unawaited(abortTrigger.whenComplete(cronetRequest.cancel));
}
cronetRequest.start();
return responseCompleter.future;
try {
cronetRequest.start();

return await responseCompleter.future;
} finally {
await disposeUpload?.call();
}
});
}
}
Expand Down
Loading