writeBlob method

  1. @override
Future<BlobWriteResult> writeBlob(
  1. BlobWriteRequest request
)
override

Write blob from a byte stream; should enforce optimistic versioning when expectedVersion is provided in the request.

Implementation

@override
Future<BlobWriteResult> writeBlob(BlobWriteRequest request) async {
  _ensureOpen();
  final id = request.id ?? _generateId();
  final payload = await _collectBytes(
    request.bytes,
    declaredLength: request.length,
  );

  if (_maxBlobBytes != null && payload.length > _maxBlobBytes) {
    throw StateError(
      'Blob too large: ${payload.length} bytes exceeds $_maxBlobBytes',
    );
  }

  if (request.checksum != null) {
    _verifyChecksum(
      payload,
      request.checksum!,
      algorithm: request.checksumAlgorithm,
    );
  }

  final now = _clock();
  final collectionMap = _storage.putIfAbsent(request.collection, () => {});
  final existing = collectionMap[id];

  DateTime createdAt;
  int version;

  if (existing == null) {
    if (request.expectedVersion != null) {
      throw StateError(
        'Expected version ${request.expectedVersion} for $id but blob is missing.',
      );
    }
    createdAt = now;
    version = 1;
  } else {
    final currentVersion = existing.descriptor.version;
    if (request.expectedVersion != null &&
        currentVersion != request.expectedVersion) {
      throw StateError(
        'Expected version ${request.expectedVersion} for $id, '
        'found $currentVersion.',
      );
    }

    final existingChecksum = existing.descriptor.checksum;
    if (existingChecksum != null &&
        request.checksum != null &&
        existingChecksum != request.checksum) {
      throw StateError(
        'Checksum mismatch for existing blob $id: stored=$existingChecksum new=${request.checksum}',
      );
    }

    if (request.checksum != null &&
        request.checksum == existingChecksum &&
        existing.descriptor.length == payload.length &&
        request.contentType == existing.descriptor.contentType) {
      return BlobWriteResult(descriptor: existing.descriptor);
    }

    createdAt = existing.descriptor.createdAt;
    version = currentVersion + 1;
  }

  final descriptor = BlobDescriptor(
    id: id,
    collection: request.collection,
    length: payload.length,
    version: version,
    createdAt: createdAt,
    updatedAt: now,
    contentType: request.contentType,
    checksum: request.checksum,
    metadata: Map<String, String>.from(request.metadata),
  );

  collectionMap[id] = _BlobEntry(descriptor: descriptor, payload: payload);

  return BlobWriteResult(descriptor: descriptor);
}