streamQueryNamed method

Stream<Result<QueryResult>> streamQueryNamed(
  1. String connectionId,
  2. String sql,
  3. Map<String, Object?> namedParams, {
  4. int fetchSize = 1000,
  5. int? chunkSize,
})

Implementation

Stream<Result<QueryResult>> streamQueryNamed(
  String connectionId,
  String sql,
  Map<String, Object?> namedParams, {
  int fetchSize = 1000,
  int? chunkSize,
}) async* {
  final nativeId = state.connectionIds[connectionId];
  if (nativeId == null) {
    yield const Failure<QueryResult, OdbcError>(
      ValidationError(message: 'Invalid connection ID'),
    );
    return;
  }

  late final String cleanedSql;
  late final Uint8List paramsBuffer;
  try {
    final extract = NamedParameterParser.extract(sql);
    final positional = NamedParameterParser.toPositionalParams(
      namedParams: namedParams,
      paramNames: extract.paramNames,
    );
    cleanedSql = extract.cleanedSql;
    paramsBuffer = serializeParams(paramValuesFromObjects(positional));
  } on ParameterMissingException catch (e) {
    yield Failure<QueryResult, OdbcError>(
      ValidationError(message: e.message),
    );
    return;
  } on Exception catch (e) {
    yield Failure<QueryResult, OdbcError>(
      QueryError(message: e.toString()),
    );
    return;
  }

  final supportsParams =
      ffi.isAsync || ffi.sync.native.supportsStreamStartParams;
  if (!supportsParams) {
    yield await query.executeQueryNamed(connectionId, sql, namedParams);
    return;
  }

  final opts = state.optionsFor(connectionId);
  final effectiveChunk = resolveStreamChunkSizeBytes(
    chunkSize: chunkSize,
    options: opts,
  );
  final maxBytes = opts?.maxResultBufferBytes;
  final queryTimeout = opts?.queryTimeout;
  final lazyStrings = opts?.lazyStrings ?? false;

  Stream<Result<QueryResult>> createSource() async* {
    try {
      await for (final chunk in streamNativeQueryWithFallback(
        nativeId,
        cleanedSql,
        maxBufferBytes: maxBytes,
        lazyStrings: lazyStrings,
        paramsBuffer: paramsBuffer,
        fetchSize: fetchSize,
        chunkSize: effectiveChunk,
      )) {
        yield Success(parser.toQueryResult(chunk));
      }
    } on Exception catch (e) {
      yield await _errors.streamingFailureFromException(e);
    }
  }

  yield* streamWithQueryTimeout(
    source: createSource(),
    queryTimeout: queryTimeout,
    onTimeoutItem: const Failure<QueryResult, OdbcError>(
      QueryError(message: odbcQueryTimedOutMessage),
    ),
  );
}