Skip to content

CarpDataStreamService

A CARP Core DataStreamService that talks to CAWS.

Uploads collected Measurements to CAWS as DataStreamBatches, and reads them back. Data streams are opened, closed and removed on the server, not by the client.

Key points:

Apps rarely call it directly: the carp_backend package uploads data through this service when the protocol uses a CAWS data endpoint.

Inheritance
Implemented types

Constructors

CarpDataStreamService

CarpDataStreamService()

Returns the singleton default instance of the CarpDataStreamService. Before this instance can be used, it must be configured using the configure method.

Constants

DATA_STREAM_ENDPOINT_NAME

String const DATA_STREAM_ENDPOINT_NAME

The default (uncompressed) RPC endpoint name.

DATA_STREAM_QUERY_BY_TIME_ENDPOINT_NAME

String const DATA_STREAM_QUERY_BY_TIME_ENDPOINT_NAME

The endpoint name used by getDataStreamBatchesByTime.

DATA_STREAM_ZIP_ENDPOINT_NAME

String const DATA_STREAM_ZIP_ENDPOINT_NAME

The endpoint name for gzip-compressed uploads.

Properties

rpcEndpointName

String get rpcEndpointName
override

The name of this service's RPC endpoint at CAWS, like deployment-service.

Methods

appendToDataStreams

Future<void> appendToDataStreams(
  1. String studyDeploymentId,
  2. List<DataStreamBatch> batch, {
  3. bool compress = true,
})
override

Appends a batch of data to the data streams of the study deployment with studyDeploymentId.

If compress is true (default), the JSON payload is gzipped before upload.

closeDataStreams

Future<void> closeDataStreams(
  1. List<String> studyDeploymentIds
)
override

Stop accepting incoming data for all data streams for each of the studyDeploymentIds.

Throws IllegalArgumentException when no data streams were ever opened for any of the studyDeploymentIds.

dataStream

DataStreamReference dataStream(
  1. String studyDeploymentId
)

Gets a DataStreamReference for a studyDeploymentId.

getDataStream

Future<List<DataStreamBatch>> getDataStream(
  1. DataStreamId dataStream,
  2. int fromSequenceId, [
  3. int? toSequenceIdInclusive
])
override

Retrieve all data points in dataStream that fall within the inclusive range defined by fromSequenceId and toSequenceIdInclusive. If toSequenceIdInclusive is null, all data points starting fromSequenceId are returned.

In case no data for dataStream is stored in this repository, or is available for the specified range, an empty list is returned.

Throws IllegalArgumentException if:

  • dataStream has never been opened
  • the dataStream does not exist (i.e, that the combination of dataStream.deviceRoleName and dataStream.dataType is correct for the protocol used in the dataStream.studyDeploymentId deployment.)
  • fromSequenceId is negative or toSequenceIdInclusive is smaller than fromSequenceId

getDataStreamBatchesByTime

Future<List<DataStreamBatch>> getDataStreamBatchesByTime(
  1. DataStreamId dataStream,
  2. DateTime from,
  3. DateTime to
)

Query dataStream by its local update time window instead of a sequence-id range.

Returns all data points in dataStream whose local updated_at timestamp falls within the inclusive from-to window, as one DataStreamBatch per contiguous run of measurements - a new batch starts wherever the sequence was interrupted.

This is a CAWS-specific endpoint (not part of the core DataStreamService interface) and mirrors getDataStream, but is useful when the local upload time is more relevant than the sequence id (e.g., incremental sync of recently uploaded data).

openDataStreams

Future<void> openDataStreams(
  1. DataStreamsConfiguration configuration
)
override

Start accepting data for a study deployment for data streams configured in configuration.

Throws IllegalStateException when data streams for the specified study deployment have already been configured.

removeDataStreams

Future<Set<String>> removeDataStreams(
  1. List<String> studyDeploymentIds
)
override

Close data streams and remove all data for each of the studyDeploymentIds.

Returns the IDs of the study deployments for which data streams were configured. IDs for which no study deployment exists are ignored.

stream

  1. @Deprecated('Use dataStream() instead.')
DataStreamReference stream(
  1. String studyDeploymentId
)

Gets a DataStreamReference for a studyDeploymentId.