From 7909ef801278f0e19d776cc5913c10f6621483c5 Mon Sep 17 00:00:00 2001 From: "PUBNUB\\jakub.grzesiowski" Date: Tue, 7 Apr 2026 17:17:13 +0200 Subject: [PATCH 01/31] Implement DataSync CRUD --- .../DataSync/CreateEntityOperation.cs | 238 +++++++++++++++++ .../DataSync/CreateEntityParameters.cs | 44 ++++ .../PubnubApi/EndPoint/DataSync/DataSync.cs | 92 +++++++ .../DataSync/DeleteEntityOperation.cs | 175 +++++++++++++ .../DataSync/DeleteEntityParameters.cs | 19 ++ .../EndPoint/DataSync/GetEntitiesOperation.cs | 226 ++++++++++++++++ .../DataSync/GetEntitiesParameters.cs | 81 ++++++ .../EndPoint/DataSync/GetEntityOperation.cs | 180 +++++++++++++ .../EndPoint/DataSync/GetEntityParameters.cs | 15 ++ .../DataSync/PNDataSyncDeleteEntityResult.cs | 5 + .../DataSync/PNDataSyncEntityResult.cs | 20 ++ .../EndPoint/DataSync/PatchEntityOperation.cs | 245 ++++++++++++++++++ .../DataSync/PatchEntityParameters.cs | 66 +++++ .../DataSync/UpdateEntityOperation.cs | 222 ++++++++++++++++ .../DataSync/UpdateEntityParameters.cs | 38 +++ src/Api/PubnubApi/Enum/PNOperationType.cs | 9 +- .../DeserializeToInternalObjectUtility.cs | 19 ++ ...DataSyncEntitiesListResultJsonDataParse.cs | 134 ++++++++++ .../PNDataSyncEntityResultJsonDataParse.cs | 59 +++++ src/Api/PubnubApi/Model/Consumer/PNStatus.cs | 60 +++++ .../PNDataSyncDeleteEntityResultExt.cs | 19 ++ .../PNDataSyncEntitiesListResultExt.cs | 19 ++ .../DataSync/PNDataSyncEntityResultExt.cs | 19 ++ src/Api/PubnubApi/Pubnub.cs | 3 + src/Api/PubnubApi/PubnubCoreBase.cs | 14 + .../PubnubApi/Transport/HttpClientService.cs | 44 +++- src/Api/PubnubApiPCL/PubnubApiPCL.csproj | 84 ++++++ src/Api/PubnubApiUWP/PubnubApiUWP.csproj | 83 ++++++ src/Api/PubnubApiUnity/PubnubApiUnity.csproj | 84 ++++++ .../PubnubApiPCL.Tests.csproj | 3 + 30 files changed, 2308 insertions(+), 11 deletions(-) create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/CreateEntityOperation.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/CreateEntityParameters.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/DataSync.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/DeleteEntityOperation.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/DeleteEntityParameters.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/GetEntitiesOperation.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/GetEntitiesParameters.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/GetEntityOperation.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/GetEntityParameters.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/PNDataSyncDeleteEntityResult.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/PNDataSyncEntityResult.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/PatchEntityOperation.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/PatchEntityParameters.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/UpdateEntityOperation.cs create mode 100644 src/Api/PubnubApi/EndPoint/DataSync/UpdateEntityParameters.cs create mode 100644 src/Api/PubnubApi/JsonDataParse/PNDataSyncEntitiesListResultJsonDataParse.cs create mode 100644 src/Api/PubnubApi/JsonDataParse/PNDataSyncEntityResultJsonDataParse.cs create mode 100644 src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncDeleteEntityResultExt.cs create mode 100644 src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncEntitiesListResultExt.cs create mode 100644 src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncEntityResultExt.cs diff --git a/src/Api/PubnubApi/EndPoint/DataSync/CreateEntityOperation.cs b/src/Api/PubnubApi/EndPoint/DataSync/CreateEntityOperation.cs new file mode 100644 index 000000000..d76ce3c59 --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/CreateEntityOperation.cs @@ -0,0 +1,238 @@ +using System; +using System.Collections.Generic; +using System.Text; +using System.Threading.Tasks; + +namespace PubnubApi.EndPoint +{ + public class CreateEntityOperation : PubnubCoreBase + { + private readonly PNConfiguration config; + private readonly IJsonPluggableLibrary jsonLibrary; + private readonly IPubnubUnitTest unit; + private readonly CreateEntityParameters parameters; + + private PNCallback savedCallback; + + private const PNOperationType OperationType = PNOperationType.PNDataSyncCreateEntity; + + public CreateEntityOperation(PNConfiguration pubnubConfig, IJsonPluggableLibrary jsonPluggableLibrary, + IPubnubUnitTest pubnubUnit, TokenManager tokenManager, Pubnub instance, + CreateEntityParameters parameters) : base(pubnubConfig, + jsonPluggableLibrary, pubnubUnit, tokenManager, instance) + { + config = pubnubConfig; + jsonLibrary = jsonPluggableLibrary; + unit = pubnubUnit; + this.parameters = parameters ?? throw new ArgumentNullException(nameof(parameters)); + } + + public void Execute(PNCallback callback) + { + if (callback == null) + { + throw new ArgumentException("Missing userCallback"); + } + + savedCallback = callback; + ExecuteAsync().ContinueWith(t => + { + if (t.IsFaulted) + { + var status = new PNStatus + { + Error = true, + ErrorData = new PNErrorData(t.Exception?.Message, t.Exception) + }; + callback.OnResponse(default, status); + } + else + { + var pnResult = t.Result; + callback.OnResponse(pnResult.Result, pnResult.Status); + } + }); + } + + public async Task> ExecuteAsync() + { + logger?.Trace($"{GetType().Name} ExecuteAsync invoked."); + return await CreateEntityAsync().ConfigureAwait(false); + } + + internal void Retry() + { + if (savedCallback != null) + { + Execute(savedCallback); + } + } + + private async Task> CreateEntityAsync() + { + var returnValue = new PNResult(); + + if (string.IsNullOrEmpty(parameters.EntityClass) || string.IsNullOrEmpty(parameters.EntityClass.Trim())) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Missing EntityClass", + new ArgumentException("Missing EntityClass")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + if (parameters.EntityClassVersion < 1) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("EntityClassVersion must be >= 1", + new ArgumentException("EntityClassVersion must be >= 1")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + if (string.IsNullOrEmpty(parameters.IdempotencyKey) || string.IsNullOrEmpty(parameters.IdempotencyKey.Trim())) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Missing IdempotencyKey", + new ArgumentException("Missing IdempotencyKey")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + if (string.IsNullOrEmpty(config.SubscribeKey) || string.IsNullOrEmpty(config.SubscribeKey.Trim()) || + config.SubscribeKey.Length <= 0) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Invalid Subscribe key", + new ArgumentException("Invalid Subscribe key")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + logger?.Trace($"{GetType().Name} parameter validated."); + var requestState = new RequestState + { + ResponseType = OperationType, + Reconnect = false, + EndPointOperation = this, + UsePostMethod = true + }; + + var requestParameter = CreateRequestParameter(); + Tuple JsonAndStatusTuple; + var transportRequest = PubnubInstance.transportMiddleware.PreapareTransportRequest( + requestParameter, OperationType); + var transportResponse = await PubnubInstance.transportMiddleware.Send(transportRequest) + .ConfigureAwait(false); + if (transportResponse.Error == null) + { + var responseString = Encoding.UTF8.GetString(transportResponse.Content); + var errorStatus = GetStatusIfError(requestState, responseString); + //TODO: why 201 and not 200 from constants? + if (errorStatus == null && transportResponse.StatusCode == 201) + { + requestState.GotJsonResponse = true; + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + requestState.ResponseType, PNStatusCategory.PNAcknowledgmentCategory, requestState, 201, + null); + JsonAndStatusTuple = new Tuple(responseString, status); + } + else + { + JsonAndStatusTuple = new Tuple(string.Empty, errorStatus); + } + + returnValue.Status = JsonAndStatusTuple.Item2; + var json = JsonAndStatusTuple.Item1; + if (!string.IsNullOrEmpty(json)) + { + var resultList = ProcessJsonResponse(requestState, json); + var responseBuilder = new ResponseBuilder(config, jsonLibrary); + var responseResult = + responseBuilder.JsonToObject(resultList, true); + if (responseResult != null) + { + returnValue.Result = responseResult; + } + } + } + else + { + var statusCode = PNStatusCodeHelper.GetHttpStatusCode(transportResponse.Error.Message); + var category = + PNStatusCategoryHelper.GetPNStatusCategory(statusCode, transportResponse.Error.Message); + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + OperationType, category, requestState, statusCode, + new PNException(transportResponse.Error.Message, transportResponse.Error)); + returnValue.Status = status; + } + + logger?.Trace($"{GetType().Name} request finished with status code {returnValue.Status.StatusCode}"); + return returnValue; + } + + private RequestParameter CreateRequestParameter() + { + var dataProperties = new Dictionary(); + + if (!string.IsNullOrEmpty(parameters.Id)) + { + dataProperties.Add("id", parameters.Id); + } + + dataProperties.Add("entityClass", parameters.EntityClass); + dataProperties.Add("entityClassVersion", parameters.EntityClassVersion); + + if (!string.IsNullOrEmpty(parameters.Status)) + { + dataProperties.Add("status", parameters.Status); + } + + if (parameters.Payload != null) + { + dataProperties.Add("payload", parameters.Payload); + } + + var requestEnvelope = new Dictionary + { + { "data", dataProperties } + }; + + var postBody = jsonLibrary.SerializeToJsonString(requestEnvelope); + + var pathSegments = new List + { + //"v1", + //"datasync", + "subkeys", + config.SubscribeKey, + "entities" + }; + + var requestParameter = new RequestParameter + { + RequestType = Constants.POST, + PathSegment = pathSegments, + BodyContentString = postBody + }; + + requestParameter.Headers.Add("Content-Type", + "application/vnd.pubnub.objects.entity+json;version=1"); + requestParameter.Headers.Add("Idempotency-Key", parameters.IdempotencyKey); + + return requestParameter; + } + } +} diff --git a/src/Api/PubnubApi/EndPoint/DataSync/CreateEntityParameters.cs b/src/Api/PubnubApi/EndPoint/DataSync/CreateEntityParameters.cs new file mode 100644 index 000000000..99967ab91 --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/CreateEntityParameters.cs @@ -0,0 +1,44 @@ +using System; +using System.Collections.Generic; + +namespace PubnubApi.EndPoint +{ + /// + /// Parameters for creating a new entity via DataSync. + /// + public class CreateEntityParameters + { + /// + /// Entity identifier. Optional — if not provided, the server generates one. + /// Must be 1–255 characters if provided. + /// + public string Id { get; set; } + + /// + /// Entity class identifier (e.g., "vehicle", "order", "sensor"). + /// Required. Immutable after creation. + /// + public string EntityClass { get; set; } + + /// + /// Schema version of the entity class. Required. Must be >= 1. + /// + public int EntityClassVersion { get; set; } + + /// + /// Entity status (e.g., "active", "inactive"). 1–100 characters. + /// + public string Status { get; set; } + + /// + /// User-defined custom properties. Supports arbitrarily nested objects. + /// + public Dictionary Payload { get; set; } + + /// + /// Idempotency key (UUIDv4) to ensure the request is processed exactly once. + /// Required for POST requests. + /// + public string IdempotencyKey { get; set; } + } +} diff --git a/src/Api/PubnubApi/EndPoint/DataSync/DataSync.cs b/src/Api/PubnubApi/EndPoint/DataSync/DataSync.cs new file mode 100644 index 000000000..ca310afb2 --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/DataSync.cs @@ -0,0 +1,92 @@ +using System.Threading.Tasks; + +namespace PubnubApi.EndPoint; + +/// +/// Entrypoint for PubNub Data Sync operations +/// +public class DataSync +{ + private readonly IPubnubUnitTest unit; + private readonly Pubnub pubnub; + private readonly TokenManager tokenManager; + + public DataSync(Pubnub pubnub, IPubnubUnitTest unit, TokenManager tokenManager) + { + this.pubnub = pubnub; + this.unit = unit; + this.tokenManager = tokenManager; + } + + public async Task> CreateEntity(CreateEntityParameters parameters) + { + return await new CreateEntityOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).ExecuteAsync(); + } + + public void CreateEntity(CreateEntityParameters parameters, PNDataSyncEntityResultExt callback) + { + new CreateEntityOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).Execute(callback); + } + + public async Task> GetEntity(GetEntityParameters parameters) + { + return await new GetEntityOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).ExecuteAsync(); + } + + public void GetEntity(GetEntityParameters parameters, PNDataSyncEntityResultExt callback) + { + new GetEntityOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).Execute(callback); + } + + public async Task> GetEntities(GetEntitiesParameters parameters) + { + return await new GetEntitiesOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).ExecuteAsync(); + } + + public void GetEntities(GetEntitiesParameters parameters, PNDataSyncEntitiesListResultExt callback) + { + new GetEntitiesOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).Execute(callback); + } + + public async Task> UpdateEntity(UpdateEntityParameters parameters) + { + return await new UpdateEntityOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).ExecuteAsync(); + } + + public void UpdateEntity(UpdateEntityParameters parameters, PNDataSyncEntityResultExt callback) + { + new UpdateEntityOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).Execute(callback); + } + + public async Task> PatchEntity(PatchEntityParameters parameters) + { + return await new PatchEntityOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).ExecuteAsync(); + } + + public void PatchEntity(PatchEntityParameters parameters, PNDataSyncEntityResultExt callback) + { + new PatchEntityOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).Execute(callback); + } + + public async Task> DeleteEntity(DeleteEntityParameters parameters) + { + return await new DeleteEntityOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).ExecuteAsync(); + } + + public void DeleteEntity(DeleteEntityParameters parameters, PNDataSyncDeleteEntityResultExt callback) + { + new DeleteEntityOperation(pubnub.PNConfig, pubnub.JsonPluggableLibrary, unit, tokenManager, pubnub, + parameters).Execute(callback); + } +} \ No newline at end of file diff --git a/src/Api/PubnubApi/EndPoint/DataSync/DeleteEntityOperation.cs b/src/Api/PubnubApi/EndPoint/DataSync/DeleteEntityOperation.cs new file mode 100644 index 000000000..9a68813dd --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/DeleteEntityOperation.cs @@ -0,0 +1,175 @@ +using System; +using System.Collections.Generic; +using System.Text; +using System.Threading.Tasks; + +namespace PubnubApi.EndPoint +{ + public class DeleteEntityOperation : PubnubCoreBase + { + private readonly PNConfiguration config; + private readonly IJsonPluggableLibrary jsonLibrary; + private readonly IPubnubUnitTest unit; + private readonly DeleteEntityParameters parameters; + + private PNCallback savedCallback; + + private const PNOperationType OperationType = PNOperationType.PNDataSyncDeleteEntity; + + public DeleteEntityOperation(PNConfiguration pubnubConfig, IJsonPluggableLibrary jsonPluggableLibrary, + IPubnubUnitTest pubnubUnit, TokenManager tokenManager, Pubnub instance, + DeleteEntityParameters parameters) : base(pubnubConfig, + jsonPluggableLibrary, pubnubUnit, tokenManager, instance) + { + config = pubnubConfig; + jsonLibrary = jsonPluggableLibrary; + unit = pubnubUnit; + this.parameters = parameters ?? throw new ArgumentNullException(nameof(parameters)); + } + + public void Execute(PNCallback callback) + { + if (callback == null) + { + throw new ArgumentException("Missing userCallback"); + } + + savedCallback = callback; + ExecuteAsync().ContinueWith(t => + { + if (t.IsFaulted) + { + var status = new PNStatus + { + Error = true, + ErrorData = new PNErrorData(t.Exception?.Message, t.Exception) + }; + callback.OnResponse(default, status); + } + else + { + var pnResult = t.Result; + callback.OnResponse(pnResult.Result, pnResult.Status); + } + }); + } + + public async Task> ExecuteAsync() + { + logger?.Trace($"{GetType().Name} ExecuteAsync invoked."); + return await DeleteEntityAsync().ConfigureAwait(false); + } + + internal void Retry() + { + if (savedCallback != null) + { + Execute(savedCallback); + } + } + + private async Task> DeleteEntityAsync() + { + var returnValue = new PNResult(); + + if (string.IsNullOrEmpty(parameters.Id) || string.IsNullOrEmpty(parameters.Id.Trim())) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Missing Entity Id", + new ArgumentException("Missing Entity Id")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + if (string.IsNullOrEmpty(config.SubscribeKey) || string.IsNullOrEmpty(config.SubscribeKey.Trim()) || + config.SubscribeKey.Length <= 0) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Invalid Subscribe key", + new ArgumentException("Invalid Subscribe key")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + logger?.Trace($"{GetType().Name} parameter validated."); + var requestState = new RequestState + { + ResponseType = OperationType, + Reconnect = false, + EndPointOperation = this, + UsePostMethod = false + }; + + var requestParameter = CreateRequestParameter(); + var transportRequest = PubnubInstance.transportMiddleware.PreapareTransportRequest( + requestParameter, OperationType); + var transportResponse = await PubnubInstance.transportMiddleware.Send(transportRequest) + .ConfigureAwait(false); + if (transportResponse.Error == null) + { + var responseString = Encoding.UTF8.GetString(transportResponse.Content); + var errorStatus = GetStatusIfError(requestState, responseString); + Tuple JsonAndStatusTuple; + if (transportResponse.StatusCode == Constants.HttpRequestSuccessStatusCode) + { + logger?.Trace($"{GetType().Name} request finished with status code {transportResponse.StatusCode}"); + requestState.GotJsonResponse = true; + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + requestState.ResponseType, PNStatusCategory.PNAcknowledgmentCategory, requestState, + Constants.HttpRequestSuccessStatusCode, null); + JsonAndStatusTuple = new Tuple(responseString, status); + } + else + { + JsonAndStatusTuple = new Tuple(string.Empty, errorStatus); + } + + returnValue.Status = JsonAndStatusTuple.Item2; + returnValue.Result = new PNDataSyncDeleteEntityResult(); + } + else + { + var statusCode = PNStatusCodeHelper.GetHttpStatusCode(transportResponse.Error.Message); + var category = + PNStatusCategoryHelper.GetPNStatusCategory(statusCode, transportResponse.Error.Message); + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + OperationType, category, requestState, statusCode, + new PNException(transportResponse.Error.Message, transportResponse.Error)); + returnValue.Status = status; + } + + logger?.Trace($"{GetType().Name} request finished with status code {returnValue.Status.StatusCode}"); + return returnValue; + } + + private RequestParameter CreateRequestParameter() + { + var pathSegments = new List + { + "subkeys", + config.SubscribeKey, + "entities", + parameters.Id + }; + + var requestParameter = new RequestParameter + { + RequestType = Constants.DELETE, + PathSegment = pathSegments + }; + + if (!string.IsNullOrEmpty(parameters.IfMatch)) + { + requestParameter.Headers.Add("If-Match", parameters.IfMatch); + } + + return requestParameter; + } + } +} diff --git a/src/Api/PubnubApi/EndPoint/DataSync/DeleteEntityParameters.cs b/src/Api/PubnubApi/EndPoint/DataSync/DeleteEntityParameters.cs new file mode 100644 index 000000000..926a6025e --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/DeleteEntityParameters.cs @@ -0,0 +1,19 @@ +namespace PubnubApi.EndPoint +{ + /// + /// Parameters for deleting an entity by ID via DataSync. + /// + public class DeleteEntityParameters + { + /// + /// Entity identifier. Required. + /// + public string Id { get; set; } + + /// + /// ETag for optimistic concurrency control. If provided, the server rejects + /// the delete when the current resource version does not match (HTTP 412). + /// + public string IfMatch { get; set; } + } +} diff --git a/src/Api/PubnubApi/EndPoint/DataSync/GetEntitiesOperation.cs b/src/Api/PubnubApi/EndPoint/DataSync/GetEntitiesOperation.cs new file mode 100644 index 000000000..2a1c38161 --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/GetEntitiesOperation.cs @@ -0,0 +1,226 @@ +using System; +using System.Collections.Generic; +using System.Text; +using System.Threading.Tasks; + +namespace PubnubApi.EndPoint +{ + public class GetEntitiesOperation : PubnubCoreBase + { + private readonly PNConfiguration config; + private readonly IJsonPluggableLibrary jsonLibrary; + private readonly IPubnubUnitTest unit; + private readonly GetEntitiesParameters parameters; + + private PNCallback savedCallback; + + private const PNOperationType OperationType = PNOperationType.PNDataSyncGetEntities; + + public GetEntitiesOperation(PNConfiguration pubnubConfig, IJsonPluggableLibrary jsonPluggableLibrary, + IPubnubUnitTest pubnubUnit, TokenManager tokenManager, Pubnub instance, + GetEntitiesParameters parameters) : base(pubnubConfig, + jsonPluggableLibrary, pubnubUnit, tokenManager, instance) + { + config = pubnubConfig; + jsonLibrary = jsonPluggableLibrary; + unit = pubnubUnit; + this.parameters = parameters ?? throw new ArgumentNullException(nameof(parameters)); + } + + public void Execute(PNCallback callback) + { + if (callback == null) + { + throw new ArgumentException("Missing userCallback"); + } + + savedCallback = callback; + ExecuteAsync().ContinueWith(t => + { + if (t.IsFaulted) + { + var status = new PNStatus + { + Error = true, + ErrorData = new PNErrorData(t.Exception?.Message, t.Exception) + }; + callback.OnResponse(default, status); + } + else + { + var pnResult = t.Result; + callback.OnResponse(pnResult.Result, pnResult.Status); + } + }); + } + + public async Task> ExecuteAsync() + { + logger?.Trace($"{GetType().Name} ExecuteAsync invoked."); + return await GetEntitiesAsync().ConfigureAwait(false); + } + + internal void Retry() + { + if (savedCallback != null) + { + Execute(savedCallback); + } + } + + private async Task> GetEntitiesAsync() + { + var returnValue = new PNResult(); + + if (string.IsNullOrEmpty(parameters.EntityClass) || string.IsNullOrEmpty(parameters.EntityClass.Trim())) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Missing EntityClass", + new ArgumentException("Missing EntityClass")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + if (string.IsNullOrEmpty(config.SubscribeKey) || string.IsNullOrEmpty(config.SubscribeKey.Trim()) || + config.SubscribeKey.Length <= 0) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Invalid Subscribe key", + new ArgumentException("Invalid Subscribe key")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + logger?.Trace($"{GetType().Name} parameter validated."); + var requestState = new RequestState + { + ResponseType = OperationType, + Reconnect = false, + EndPointOperation = this, + UsePostMethod = false + }; + + var requestParameter = CreateRequestParameter(); + var transportRequest = PubnubInstance.transportMiddleware.PreapareTransportRequest( + requestParameter, OperationType); + var transportResponse = await PubnubInstance.transportMiddleware.Send(transportRequest) + .ConfigureAwait(false); + if (transportResponse.Error == null) + { + var responseString = Encoding.UTF8.GetString(transportResponse.Content); + var errorStatus = GetStatusIfError(requestState, responseString); + Tuple JsonAndStatusTuple; + if (errorStatus == null && transportResponse.StatusCode == Constants.HttpRequestSuccessStatusCode) + { + requestState.GotJsonResponse = true; + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + requestState.ResponseType, PNStatusCategory.PNAcknowledgmentCategory, requestState, + Constants.HttpRequestSuccessStatusCode, null); + JsonAndStatusTuple = new Tuple(responseString, status); + } + else + { + JsonAndStatusTuple = new Tuple(string.Empty, errorStatus); + } + + returnValue.Status = JsonAndStatusTuple.Item2; + var json = JsonAndStatusTuple.Item1; + if (!string.IsNullOrEmpty(json)) + { + var resultList = ProcessJsonResponse(requestState, json); + var responseBuilder = new ResponseBuilder(config, jsonLibrary); + var responseResult = + responseBuilder.JsonToObject(resultList, true); + if (responseResult != null) + { + returnValue.Result = responseResult; + } + } + } + else + { + var statusCode = PNStatusCodeHelper.GetHttpStatusCode(transportResponse.Error.Message); + var category = + PNStatusCategoryHelper.GetPNStatusCategory(statusCode, transportResponse.Error.Message); + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + OperationType, category, requestState, statusCode, + new PNException(transportResponse.Error.Message, transportResponse.Error)); + returnValue.Status = status; + } + + logger?.Trace($"{GetType().Name} request finished with status code {returnValue.Status.StatusCode}"); + return returnValue; + } + + private RequestParameter CreateRequestParameter() + { + var pathSegments = new List + { + "subkeys", + config.SubscribeKey, + "entities" + }; + + var requestQueryStringParams = new Dictionary(); + + requestQueryStringParams.Add("entity_class", + UriUtil.EncodeUriComponent(parameters.EntityClass, + OperationType, false, false, false)); + + if (parameters.EntityClassVersion.HasValue) + { + requestQueryStringParams.Add("entity_class_version", + parameters.EntityClassVersion.Value.ToString()); + } + + if (!string.IsNullOrEmpty(parameters.Cursor)) + { + requestQueryStringParams.Add("cursor", + UriUtil.EncodeUriComponent(parameters.Cursor, + OperationType, false, false, false)); + } + + if (parameters.Limit.HasValue) + { + requestQueryStringParams.Add("limit", + parameters.Limit.Value.ToString()); + } + + if (!string.IsNullOrEmpty(parameters.Filter)) + { + requestQueryStringParams.Add("filter", + UriUtil.EncodeUriComponent(parameters.Filter, + OperationType, false, false, false)); + } + + if (!string.IsNullOrEmpty(parameters.FilterAdvanced)) + { + requestQueryStringParams.Add("filter_advanced", + UriUtil.EncodeUriComponent(parameters.FilterAdvanced, + OperationType, false, false, false)); + } + + if (!string.IsNullOrEmpty(parameters.Sort)) + { + requestQueryStringParams.Add("sort", + UriUtil.EncodeUriComponent(parameters.Sort, + OperationType, false, false, false)); + } + + var requestParameter = new RequestParameter + { + RequestType = Constants.GET, + PathSegment = pathSegments, + Query = requestQueryStringParams + }; + + return requestParameter; + } + } +} diff --git a/src/Api/PubnubApi/EndPoint/DataSync/GetEntitiesParameters.cs b/src/Api/PubnubApi/EndPoint/DataSync/GetEntitiesParameters.cs new file mode 100644 index 000000000..840cb168a --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/GetEntitiesParameters.cs @@ -0,0 +1,81 @@ +using System.Collections.Generic; + +namespace PubnubApi.EndPoint +{ + /// + /// Parameters for listing entities via DataSync. + /// + public class GetEntitiesParameters + { + /// + /// Entity class name to filter by. Required. + /// + public string EntityClass { get; set; } + + /// + /// Schema version of the entity class. Optional — if not provided the server + /// returns entities matching the latest version. + /// + public int? EntityClassVersion { get; set; } + + /// + /// Pagination cursor returned from a previous request. + /// + public string Cursor { get; set; } + + /// + /// Maximum number of items to return per page. + /// Min 1, max 100, default 20. + /// + public int? Limit { get; set; } + + /// + /// Filter expression using AppContext Query Language (e.g., "status == 'active'"). + /// + public string Filter { get; set; } + + /// + /// Advanced filter expression supporting logical operators and nested conditions. + /// + public string FilterAdvanced { get; set; } + + /// + /// Comma-separated list of fields to sort by. Prefix with + for ascending + /// or - for descending (default). Example: "-createdAt,+id". + /// + public string Sort { get; set; } + } + + /// + /// Result returned by GetEntities (list) containing an array of entities + /// plus cursor-based pagination metadata and HATEOAS links. + /// + public class PNDataSyncEntitiesListResult + { + public List Data { get; internal set; } = new(); + public PaginationMeta Meta { get; internal set; } + public PaginationLinks Links { get; internal set; } + } + + /// + /// Cursor-based pagination metadata returned in list responses. + /// + public class PaginationMeta + { + public string NextCursor { get; internal set; } + public string PrevCursor { get; internal set; } + public bool HasNext { get; internal set; } + public bool HasPrev { get; internal set; } + public int? Limit { get; internal set; } + } + + /// + /// HATEOAS navigation links returned in list responses. + /// + public class PaginationLinks + { + public string Self { get; internal set; } + public string Next { get; internal set; } + public string Prev { get; internal set; } + } +} diff --git a/src/Api/PubnubApi/EndPoint/DataSync/GetEntityOperation.cs b/src/Api/PubnubApi/EndPoint/DataSync/GetEntityOperation.cs new file mode 100644 index 000000000..f994e2180 --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/GetEntityOperation.cs @@ -0,0 +1,180 @@ +using System; +using System.Collections.Generic; +using System.Text; +using System.Threading.Tasks; + +namespace PubnubApi.EndPoint +{ + public class GetEntityOperation : PubnubCoreBase + { + private readonly PNConfiguration config; + private readonly IJsonPluggableLibrary jsonLibrary; + private readonly IPubnubUnitTest unit; + private readonly GetEntityParameters parameters; + + private PNCallback savedCallback; + + private const PNOperationType OperationType = PNOperationType.PNDataSyncGetEntity; + + public GetEntityOperation(PNConfiguration pubnubConfig, IJsonPluggableLibrary jsonPluggableLibrary, + IPubnubUnitTest pubnubUnit, TokenManager tokenManager, Pubnub instance, + GetEntityParameters parameters) : base(pubnubConfig, + jsonPluggableLibrary, pubnubUnit, tokenManager, instance) + { + config = pubnubConfig; + jsonLibrary = jsonPluggableLibrary; + unit = pubnubUnit; + this.parameters = parameters ?? throw new ArgumentNullException(nameof(parameters)); + } + + public void Execute(PNCallback callback) + { + if (callback == null) + { + throw new ArgumentException("Missing userCallback"); + } + + savedCallback = callback; + ExecuteAsync().ContinueWith(t => + { + if (t.IsFaulted) + { + var status = new PNStatus + { + Error = true, + ErrorData = new PNErrorData(t.Exception?.Message, t.Exception) + }; + callback.OnResponse(default, status); + } + else + { + var pnResult = t.Result; + callback.OnResponse(pnResult.Result, pnResult.Status); + } + }); + } + + public async Task> ExecuteAsync() + { + logger?.Trace($"{GetType().Name} ExecuteAsync invoked."); + return await GetEntityAsync().ConfigureAwait(false); + } + + internal void Retry() + { + if (savedCallback != null) + { + Execute(savedCallback); + } + } + + private async Task> GetEntityAsync() + { + var returnValue = new PNResult(); + + if (string.IsNullOrEmpty(parameters.Id) || string.IsNullOrEmpty(parameters.Id.Trim())) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Missing Entity Id", + new ArgumentException("Missing Entity Id")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + if (string.IsNullOrEmpty(config.SubscribeKey) || string.IsNullOrEmpty(config.SubscribeKey.Trim()) || + config.SubscribeKey.Length <= 0) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Invalid Subscribe key", + new ArgumentException("Invalid Subscribe key")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + logger?.Trace($"{GetType().Name} parameter validated."); + var requestState = new RequestState + { + ResponseType = OperationType, + Reconnect = false, + EndPointOperation = this, + UsePostMethod = false + }; + + var requestParameter = CreateRequestParameter(); + var transportRequest = PubnubInstance.transportMiddleware.PreapareTransportRequest( + requestParameter, OperationType); + var transportResponse = await PubnubInstance.transportMiddleware.Send(transportRequest) + .ConfigureAwait(false); + if (transportResponse.Error == null) + { + var responseString = Encoding.UTF8.GetString(transportResponse.Content); + var errorStatus = GetStatusIfError(requestState, responseString); + Tuple JsonAndStatusTuple; + if (errorStatus == null && transportResponse.StatusCode == Constants.HttpRequestSuccessStatusCode) + { + requestState.GotJsonResponse = true; + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + requestState.ResponseType, PNStatusCategory.PNAcknowledgmentCategory, requestState, + Constants.HttpRequestSuccessStatusCode, null); + JsonAndStatusTuple = new Tuple(responseString, status); + } + else + { + JsonAndStatusTuple = new Tuple(string.Empty, errorStatus); + } + + returnValue.Status = JsonAndStatusTuple.Item2; + var json = JsonAndStatusTuple.Item1; + if (!string.IsNullOrEmpty(json)) + { + var resultList = ProcessJsonResponse(requestState, json); + var responseBuilder = new ResponseBuilder(config, jsonLibrary); + var responseResult = + responseBuilder.JsonToObject(resultList, true); + if (responseResult != null) + { + returnValue.Result = responseResult; + } + } + } + else + { + var statusCode = PNStatusCodeHelper.GetHttpStatusCode(transportResponse.Error.Message); + var category = + PNStatusCategoryHelper.GetPNStatusCategory(statusCode, transportResponse.Error.Message); + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + OperationType, category, requestState, statusCode, + new PNException(transportResponse.Error.Message, transportResponse.Error)); + returnValue.Status = status; + } + + logger?.Trace($"{GetType().Name} request finished with status code {returnValue.Status.StatusCode}"); + return returnValue; + } + + private RequestParameter CreateRequestParameter() + { + var pathSegments = new List + { + "subkeys", + config.SubscribeKey, + "entities", + parameters.Id + }; + + var requestParameter = new RequestParameter + { + RequestType = Constants.GET, + PathSegment = pathSegments + }; + + return requestParameter; + } + } +} diff --git a/src/Api/PubnubApi/EndPoint/DataSync/GetEntityParameters.cs b/src/Api/PubnubApi/EndPoint/DataSync/GetEntityParameters.cs new file mode 100644 index 000000000..be6c079f8 --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/GetEntityParameters.cs @@ -0,0 +1,15 @@ +using System.Collections.Generic; + +namespace PubnubApi.EndPoint +{ + /// + /// Parameters for retrieving an entity by ID via DataSync. + /// + public class GetEntityParameters + { + /// + /// Entity identifier. Required. + /// + public string Id { get; set; } + } +} diff --git a/src/Api/PubnubApi/EndPoint/DataSync/PNDataSyncDeleteEntityResult.cs b/src/Api/PubnubApi/EndPoint/DataSync/PNDataSyncDeleteEntityResult.cs new file mode 100644 index 000000000..49705f0eb --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/PNDataSyncDeleteEntityResult.cs @@ -0,0 +1,5 @@ +namespace PubnubApi.EndPoint; + +public class PNDataSyncDeleteEntityResult +{ +} \ No newline at end of file diff --git a/src/Api/PubnubApi/EndPoint/DataSync/PNDataSyncEntityResult.cs b/src/Api/PubnubApi/EndPoint/DataSync/PNDataSyncEntityResult.cs new file mode 100644 index 000000000..8a2b6bacf --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/PNDataSyncEntityResult.cs @@ -0,0 +1,20 @@ +using System.Collections.Generic; + +namespace PubnubApi.EndPoint; + +/// +/// Result returned by a single-entity operation (Create, Get, Update, Patch). +/// The server response envelope is deserialized into this type. +/// +public class PNDataSyncEntityResult +{ + public string Id { get; internal set; } + public string EntityClass { get; internal set; } + public int EntityClassVersion { get; internal set; } + public string Status { get; internal set; } + public Dictionary Payload { get; internal set; } + public string CreatedAt { get; internal set; } + public string UpdatedAt { get; internal set; } + public string ETag { get; internal set; } + public string ExpiresAt { get; internal set; } +} \ No newline at end of file diff --git a/src/Api/PubnubApi/EndPoint/DataSync/PatchEntityOperation.cs b/src/Api/PubnubApi/EndPoint/DataSync/PatchEntityOperation.cs new file mode 100644 index 000000000..2254dd232 --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/PatchEntityOperation.cs @@ -0,0 +1,245 @@ +using System; +using System.Collections.Generic; +using System.Text; +using System.Threading.Tasks; + +namespace PubnubApi.EndPoint +{ + public class PatchEntityOperation : PubnubCoreBase + { + private readonly PNConfiguration config; + private readonly IJsonPluggableLibrary jsonLibrary; + private readonly IPubnubUnitTest unit; + private readonly PatchEntityParameters parameters; + + private PNCallback savedCallback; + + private const PNOperationType OperationType = PNOperationType.PNDataSyncPatchEntity; + + private static readonly HashSet OpsRequiringValue = + new () { JsonPatchOperationType.Add, JsonPatchOperationType.Replace, JsonPatchOperationType.Test }; + + private static readonly HashSet OpsRequiringFrom = + new () { JsonPatchOperationType.Move, JsonPatchOperationType.Copy }; + + public PatchEntityOperation(PNConfiguration pubnubConfig, IJsonPluggableLibrary jsonPluggableLibrary, + IPubnubUnitTest pubnubUnit, TokenManager tokenManager, Pubnub instance, + PatchEntityParameters parameters) : base(pubnubConfig, + jsonPluggableLibrary, pubnubUnit, tokenManager, instance) + { + config = pubnubConfig; + jsonLibrary = jsonPluggableLibrary; + unit = pubnubUnit; + this.parameters = parameters ?? throw new ArgumentNullException(nameof(parameters)); + } + + public void Execute(PNCallback callback) + { + if (callback == null) + { + throw new ArgumentException("Missing userCallback"); + } + + savedCallback = callback; + ExecuteAsync().ContinueWith(t => + { + if (t.IsFaulted) + { + var status = new PNStatus + { + Error = true, + ErrorData = new PNErrorData(t.Exception?.Message, t.Exception) + }; + callback.OnResponse(default, status); + } + else + { + var pnResult = t.Result; + callback.OnResponse(pnResult.Result, pnResult.Status); + } + }); + } + + public async Task> ExecuteAsync() + { + logger?.Trace($"{GetType().Name} ExecuteAsync invoked."); + return await PatchEntityAsync().ConfigureAwait(false); + } + + internal void Retry() + { + if (savedCallback != null) + { + Execute(savedCallback); + } + } + + private async Task> PatchEntityAsync() + { + var returnValue = new PNResult(); + + if (string.IsNullOrEmpty(parameters.Id) || string.IsNullOrEmpty(parameters.Id.Trim())) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Missing Entity Id", + new ArgumentException("Missing Entity Id")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + if (parameters.Operations == null || parameters.Operations.Count == 0) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Operations list must contain at least one operation", + new ArgumentException("Operations list must contain at least one operation")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + if (string.IsNullOrEmpty(parameters.IdempotencyKey) || + string.IsNullOrEmpty(parameters.IdempotencyKey.Trim())) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Missing IdempotencyKey", + new ArgumentException("Missing IdempotencyKey")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + if (string.IsNullOrEmpty(config.SubscribeKey) || string.IsNullOrEmpty(config.SubscribeKey.Trim()) || + config.SubscribeKey.Length <= 0) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Invalid Subscribe key", + new ArgumentException("Invalid Subscribe key")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + logger?.Trace($"{GetType().Name} parameter validated."); + var requestState = new RequestState + { + ResponseType = OperationType, + Reconnect = false, + EndPointOperation = this, + UsePostMethod = false + }; + + var requestParameter = CreateRequestParameter(); + Tuple JsonAndStatusTuple; + var transportRequest = PubnubInstance.transportMiddleware.PreapareTransportRequest( + requestParameter, OperationType); + var transportResponse = await PubnubInstance.transportMiddleware.Send(transportRequest) + .ConfigureAwait(false); + if (transportResponse.Error == null) + { + var responseString = Encoding.UTF8.GetString(transportResponse.Content); + var errorStatus = GetStatusIfError(requestState, responseString); + if (errorStatus == null && transportResponse.StatusCode == Constants.HttpRequestSuccessStatusCode) + { + requestState.GotJsonResponse = true; + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + requestState.ResponseType, PNStatusCategory.PNAcknowledgmentCategory, requestState, + Constants.HttpRequestSuccessStatusCode, null); + JsonAndStatusTuple = new Tuple(responseString, status); + } + else + { + JsonAndStatusTuple = new Tuple(string.Empty, errorStatus); + } + + returnValue.Status = JsonAndStatusTuple.Item2; + var json = JsonAndStatusTuple.Item1; + if (!string.IsNullOrEmpty(json)) + { + var resultList = ProcessJsonResponse(requestState, json); + var responseBuilder = new ResponseBuilder(config, jsonLibrary); + var responseResult = + responseBuilder.JsonToObject(resultList, true); + if (responseResult != null) + { + returnValue.Result = responseResult; + } + } + } + else + { + var statusCode = PNStatusCodeHelper.GetHttpStatusCode(transportResponse.Error.Message); + var category = + PNStatusCategoryHelper.GetPNStatusCategory(statusCode, transportResponse.Error.Message); + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + OperationType, category, requestState, statusCode, + new PNException(transportResponse.Error.Message, transportResponse.Error)); + returnValue.Status = status; + } + + logger?.Trace($"{GetType().Name} request finished with status code {returnValue.Status.StatusCode}"); + return returnValue; + } + + private RequestParameter CreateRequestParameter() + { + var patchArray = new List>(); + + foreach (var op in parameters.Operations) + { + var entry = new Dictionary + { + { "op", op.Op.ToString().ToLowerInvariant() }, + { "path", op.Path } + }; + + if (OpsRequiringValue.Contains(op.Op)) + { + entry["value"] = op.Value; + } + + if (OpsRequiringFrom.Contains(op.Op) && !string.IsNullOrEmpty(op.From)) + { + entry["from"] = op.From; + } + + patchArray.Add(entry); + } + + var patchBody = jsonLibrary.SerializeToJsonString(patchArray); + + var pathSegments = new List + { + "subkeys", + config.SubscribeKey, + "entities", + parameters.Id + }; + + var requestParameter = new RequestParameter + { + RequestType = Constants.PATCH, + PathSegment = pathSegments, + BodyContentString = patchBody + }; + + requestParameter.Headers.Add("Content-Type", "application/json-patch+json"); + requestParameter.Headers.Add("Idempotency-Key", parameters.IdempotencyKey); + + if (!string.IsNullOrEmpty(parameters.IfMatch)) + { + requestParameter.Headers.Add("If-Match", parameters.IfMatch); + } + + return requestParameter; + } + } +} diff --git a/src/Api/PubnubApi/EndPoint/DataSync/PatchEntityParameters.cs b/src/Api/PubnubApi/EndPoint/DataSync/PatchEntityParameters.cs new file mode 100644 index 000000000..0f687476a --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/PatchEntityParameters.cs @@ -0,0 +1,66 @@ +using System.Collections.Generic; + +namespace PubnubApi.EndPoint +{ + public class PatchEntityParameters + { + /// + /// Entity identifier. Required. + /// + public string Id { get; set; } + + /// + /// List of JSON Patch operations (RFC 6902) to apply. Required; must contain at least one operation. + /// + public List Operations { get; set; } + + /// + /// ETag for optimistic concurrency control. If provided, the server rejects + /// the patch when the current resource version does not match (HTTP 412). + /// + public string IfMatch { get; set; } + + /// + /// Idempotency key (UUIDv4) to ensure the request is processed exactly once. + /// Required for PATCH requests. + /// + public string IdempotencyKey { get; set; } + } + + public enum JsonPatchOperationType + { + Add, + Remove, + Replace, + Move, + Copy, + Test + } + + /// + /// A single JSON Patch operation as defined by RFC 6902. + /// + public class JsonPatchOperation + { + /// + /// The operation to perform: "add", "remove", "replace", "move", "copy", or "test". + /// + public JsonPatchOperationType Op { get; set; } + + /// + /// JSON Pointer (RFC 6901) to the target location. + /// + public string Path { get; set; } + + /// + /// The value to apply. Required for "add", "replace", and "test" operations. + /// Can be any JSON-serializable value including null. + /// + public object Value { get; set; } + + /// + /// Source location (JSON Pointer). Required for "move" and "copy" operations. + /// + public string From { get; set; } + } +} diff --git a/src/Api/PubnubApi/EndPoint/DataSync/UpdateEntityOperation.cs b/src/Api/PubnubApi/EndPoint/DataSync/UpdateEntityOperation.cs new file mode 100644 index 000000000..f7ca505a9 --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/UpdateEntityOperation.cs @@ -0,0 +1,222 @@ +using System; +using System.Collections.Generic; +using System.Text; +using System.Threading.Tasks; + +namespace PubnubApi.EndPoint +{ + public class UpdateEntityOperation : PubnubCoreBase + { + private readonly PNConfiguration config; + private readonly IJsonPluggableLibrary jsonLibrary; + private readonly IPubnubUnitTest unit; + private readonly UpdateEntityParameters parameters; + + private PNCallback savedCallback; + + private const PNOperationType OperationType = PNOperationType.PNDataSyncUpdateEntity; + + public UpdateEntityOperation(PNConfiguration pubnubConfig, IJsonPluggableLibrary jsonPluggableLibrary, + IPubnubUnitTest pubnubUnit, TokenManager tokenManager, Pubnub instance, + UpdateEntityParameters parameters) : base(pubnubConfig, + jsonPluggableLibrary, pubnubUnit, tokenManager, instance) + { + config = pubnubConfig; + jsonLibrary = jsonPluggableLibrary; + unit = pubnubUnit; + this.parameters = parameters ?? throw new ArgumentNullException(nameof(parameters)); + } + + public void Execute(PNCallback callback) + { + if (callback == null) + { + throw new ArgumentException("Missing userCallback"); + } + + savedCallback = callback; + ExecuteAsync().ContinueWith(t => + { + if (t.IsFaulted) + { + var status = new PNStatus + { + Error = true, + ErrorData = new PNErrorData(t.Exception?.Message, t.Exception) + }; + callback.OnResponse(default, status); + } + else + { + var pnResult = t.Result; + callback.OnResponse(pnResult.Result, pnResult.Status); + } + }); + } + + public async Task> ExecuteAsync() + { + logger?.Trace($"{GetType().Name} ExecuteAsync invoked."); + return await UpdateEntityAsync().ConfigureAwait(false); + } + + internal void Retry() + { + if (savedCallback != null) + { + Execute(savedCallback); + } + } + + private async Task> UpdateEntityAsync() + { + var returnValue = new PNResult(); + + if (string.IsNullOrEmpty(parameters.Id) || string.IsNullOrEmpty(parameters.Id.Trim())) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Missing Entity Id", + new ArgumentException("Missing Entity Id")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + if (parameters.EntityClassVersion < 1) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("EntityClassVersion must be >= 1", + new ArgumentException("EntityClassVersion must be >= 1")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + if (string.IsNullOrEmpty(config.SubscribeKey) || string.IsNullOrEmpty(config.SubscribeKey.Trim()) || + config.SubscribeKey.Length <= 0) + { + var errStatus = new PNStatus + { + Error = true, + ErrorData = new PNErrorData("Invalid Subscribe key", + new ArgumentException("Invalid Subscribe key")) + }; + returnValue.Status = errStatus; + return returnValue; + } + + logger?.Trace($"{GetType().Name} parameter validated."); + var requestState = new RequestState + { + ResponseType = OperationType, + Reconnect = false, + EndPointOperation = this, + UsePostMethod = false + }; + + var requestParameter = CreateRequestParameter(); + Tuple JsonAndStatusTuple; + var transportRequest = PubnubInstance.transportMiddleware.PreapareTransportRequest( + requestParameter, OperationType); + var transportResponse = await PubnubInstance.transportMiddleware.Send(transportRequest) + .ConfigureAwait(false); + if (transportResponse.Error == null) + { + var responseString = Encoding.UTF8.GetString(transportResponse.Content); + var errorStatus = GetStatusIfError(requestState, responseString); + if (errorStatus == null && transportResponse.StatusCode == Constants.HttpRequestSuccessStatusCode) + { + requestState.GotJsonResponse = true; + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + requestState.ResponseType, PNStatusCategory.PNAcknowledgmentCategory, requestState, + Constants.HttpRequestSuccessStatusCode, null); + JsonAndStatusTuple = new Tuple(responseString, status); + } + else + { + JsonAndStatusTuple = new Tuple(string.Empty, errorStatus); + } + + returnValue.Status = JsonAndStatusTuple.Item2; + var json = JsonAndStatusTuple.Item1; + if (!string.IsNullOrEmpty(json)) + { + var resultList = ProcessJsonResponse(requestState, json); + var responseBuilder = new ResponseBuilder(config, jsonLibrary); + var responseResult = + responseBuilder.JsonToObject(resultList, true); + if (responseResult != null) + { + returnValue.Result = responseResult; + } + } + } + else + { + var statusCode = PNStatusCodeHelper.GetHttpStatusCode(transportResponse.Error.Message); + var category = + PNStatusCategoryHelper.GetPNStatusCategory(statusCode, transportResponse.Error.Message); + var status = new StatusBuilder(config, jsonLibrary).CreateStatusResponse( + OperationType, category, requestState, statusCode, + new PNException(transportResponse.Error.Message, transportResponse.Error)); + returnValue.Status = status; + } + + logger?.Trace($"{GetType().Name} request finished with status code {returnValue.Status.StatusCode}"); + return returnValue; + } + + private RequestParameter CreateRequestParameter() + { + var dataProperties = new Dictionary(); + + dataProperties.Add("entityClassVersion", parameters.EntityClassVersion); + + if (!string.IsNullOrEmpty(parameters.Status)) + { + dataProperties.Add("status", parameters.Status); + } + + if (parameters.Payload != null) + { + dataProperties.Add("payload", parameters.Payload); + } + + var requestEnvelope = new Dictionary + { + { "data", dataProperties } + }; + + var putBody = jsonLibrary.SerializeToJsonString(requestEnvelope); + + var pathSegments = new List + { + "subkeys", + config.SubscribeKey, + "entities", + parameters.Id + }; + + var requestParameter = new RequestParameter + { + RequestType = Constants.PUT, + PathSegment = pathSegments, + BodyContentString = putBody + }; + + requestParameter.Headers.Add("Content-Type", + "application/vnd.pubnub.objects.entity+json;version=1"); + + if (!string.IsNullOrEmpty(parameters.IfMatch)) + { + requestParameter.Headers.Add("If-Match", parameters.IfMatch); + } + + return requestParameter; + } + } +} diff --git a/src/Api/PubnubApi/EndPoint/DataSync/UpdateEntityParameters.cs b/src/Api/PubnubApi/EndPoint/DataSync/UpdateEntityParameters.cs new file mode 100644 index 000000000..30a5c3f5f --- /dev/null +++ b/src/Api/PubnubApi/EndPoint/DataSync/UpdateEntityParameters.cs @@ -0,0 +1,38 @@ +using System.Collections.Generic; + +namespace PubnubApi.EndPoint +{ + /// + /// Parameters for updating an entity with complete resource replacement (PUT) via DataSync. + /// + public class UpdateEntityParameters + { + /// + /// Entity identifier. Required. + /// + public string Id { get; set; } + + /// + /// Schema version of the entity class. Required. Must be >= 1. + /// Note: entityClass is immutable and cannot be changed after creation. + /// + public int EntityClassVersion { get; set; } + + /// + /// Entity status (e.g., "active", "inactive"). 1–100 characters. + /// + public string Status { get; set; } + + /// + /// User-defined custom properties. Supports arbitrarily nested objects. + /// Replaces the entire payload — omitted fields are removed. + /// + public Dictionary Payload { get; set; } + + /// + /// ETag for optimistic concurrency control. If provided, the server rejects + /// the update when the current resource version does not match (HTTP 412). + /// + public string IfMatch { get; set; } + } +} diff --git a/src/Api/PubnubApi/Enum/PNOperationType.cs b/src/Api/PubnubApi/Enum/PNOperationType.cs index 106289f6b..02e047b2e 100644 --- a/src/Api/PubnubApi/Enum/PNOperationType.cs +++ b/src/Api/PubnubApi/Enum/PNOperationType.cs @@ -76,6 +76,13 @@ public enum PNOperationType PNFileUrlOperation, PNDownloadFileOperation, PNListFilesOperation, - PNDeleteFileOperation + PNDeleteFileOperation, + + PNDataSyncCreateEntity, + PNDataSyncGetEntity, + PNDataSyncGetEntities, + PNDataSyncUpdateEntity, + PNDataSyncDeleteEntity, + PNDataSyncPatchEntity } } diff --git a/src/Api/PubnubApi/JsonDataParse/DeserializeToInternalObjectUtility.cs b/src/Api/PubnubApi/JsonDataParse/DeserializeToInternalObjectUtility.cs index 1ddb4bf83..77ce6c11b 100644 --- a/src/Api/PubnubApi/JsonDataParse/DeserializeToInternalObjectUtility.cs +++ b/src/Api/PubnubApi/JsonDataParse/DeserializeToInternalObjectUtility.cs @@ -2,6 +2,7 @@ using System.Collections.Generic; using System.Globalization; using System.Linq; +using PubnubApi.EndPoint; namespace PubnubApi { @@ -969,6 +970,24 @@ public static T DeserializeToInternalObject(IJsonPluggableLibrary jsonPlug, L #endregion } + else if (typeof(T) == typeof(PNDataSyncEntityResult)) + { + #region "PNCreateEntityResult" + + PNDataSyncEntityResult result = PNDataSyncEntityResultJsonDataParse.GetObject(jsonPlug, listObject); + ret = (T)Convert.ChangeType(result, typeof(PNDataSyncEntityResult), CultureInfo.InvariantCulture); + + #endregion + } + else if (typeof(T) == typeof(PNDataSyncEntitiesListResult)) + { + #region "PNCreateEntityResult" + + PNDataSyncEntitiesListResult result = PNDataSyncEntitiesListResultJsonDataParse.GetObject(jsonPlug, listObject); + ret = (T)Convert.ChangeType(result, typeof(PNDataSyncEntitiesListResult), CultureInfo.InvariantCulture); + + #endregion + } else { System.Diagnostics.Debug.WriteLine("DeserializeToObject(list) => NO MATCH"); diff --git a/src/Api/PubnubApi/JsonDataParse/PNDataSyncEntitiesListResultJsonDataParse.cs b/src/Api/PubnubApi/JsonDataParse/PNDataSyncEntitiesListResultJsonDataParse.cs new file mode 100644 index 000000000..5a3d22c81 --- /dev/null +++ b/src/Api/PubnubApi/JsonDataParse/PNDataSyncEntitiesListResultJsonDataParse.cs @@ -0,0 +1,134 @@ +using System.Collections.Generic; +using System.Linq; +using PubnubApi.EndPoint; + +namespace PubnubApi; + +public class PNDataSyncEntitiesListResultJsonDataParse +{ + internal static PNDataSyncEntitiesListResult GetObject(IJsonPluggableLibrary jsonPlug, List listObject) + { + var result = new PNDataSyncEntitiesListResult(); + foreach (var rawObject in listObject) + { + var objectDictionary = jsonPlug.ConvertToDictionaryObject(rawObject); + if (objectDictionary == null || !objectDictionary.Any()) + { + continue; + } + + if (objectDictionary.TryGetValue("data", out var dataObject) && dataObject != null) + { + var dataArray = jsonPlug.ConvertToObjectArray(dataObject); + if (dataArray == null || !dataArray.Any()) + { + continue; + } + + for (int i = 0; i < dataArray.Length; i++) + { + if (dataArray[i] is not Dictionary dataEntryDictionary) + { + continue; + } + var data = new PNDataSyncEntityResult + { + Id = dataEntryDictionary.ContainsKey("id") && dataEntryDictionary["id"] != null + ? dataEntryDictionary["id"].ToString() + : null, + EntityClass = dataEntryDictionary.ContainsKey("entityClass") && dataEntryDictionary["entityClass"] != null + ? dataEntryDictionary["entityClass"].ToString() + : null, + EntityClassVersion = dataEntryDictionary.ContainsKey("entityClassVersion") && dataEntryDictionary["entityClassVersion"] != null + ? (int)(long)dataEntryDictionary["entityClassVersion"] + : 1, + Status = dataEntryDictionary.ContainsKey("status") && dataEntryDictionary["status"] != null + ? dataEntryDictionary["status"].ToString() + : null, + CreatedAt = dataEntryDictionary.ContainsKey("createdAt") && dataEntryDictionary["createdAt"] != null + ? dataEntryDictionary["createdAt"].ToString() + : null, + UpdatedAt = dataEntryDictionary.ContainsKey("updatedAt") && dataEntryDictionary["updatedAt"] != null + ? dataEntryDictionary["updatedAt"].ToString() + : null, + ETag = dataEntryDictionary.ContainsKey("eTag") && dataEntryDictionary["eTag"] != null + ? dataEntryDictionary["eTag"].ToString() + : null, + ExpiresAt = dataEntryDictionary.ContainsKey("expiresAt") && dataEntryDictionary["expiresAt"] != null + ? dataEntryDictionary["expiresAt"].ToString() + : null + }; + if (dataEntryDictionary.TryGetValue("payload", out var payloadObject) && payloadObject != null) + { + data.Payload = jsonPlug.ConvertToDictionaryObject(payloadObject); + } + result.Data.Add(data); + } + } + else if (objectDictionary.TryGetValue("meta", out var metaObject) && metaObject != null) + { + var metaDictionary = jsonPlug.ConvertToDictionaryObject(metaObject); + if (metaDictionary == null || !metaDictionary.Any()) + { + continue; + } + + var meta = new PaginationMeta(); + if (metaDictionary.TryGetValue("next_cursor", out var nextCursor) && nextCursor != null) + { + meta.NextCursor = nextCursor.ToString(); + } + + if (metaDictionary.TryGetValue("prev_cursor", out var prevCursor) && prevCursor != null) + { + meta.PrevCursor = prevCursor.ToString(); + } + + if (metaDictionary.TryGetValue("has_next", out var hasNext) && hasNext != null) + { + meta.HasNext = (bool)hasNext; + } + + if (metaDictionary.TryGetValue("has_prev", out var hasPrev) && hasPrev != null) + { + meta.HasPrev = (bool)hasPrev; + } + + if (metaDictionary.TryGetValue("limit", out var limit) && limit != null) + { + meta.Limit = (int)(long)limit; + } + + result.Meta = meta; + } + else if (objectDictionary.TryGetValue("links", out var linksObject) && linksObject != null) + { + var linksDictionary = jsonPlug.ConvertToDictionaryObject(listObject); + if (linksDictionary == null || !linksDictionary.Any()) + { + continue; + } + + var links = new PaginationLinks(); + if (linksDictionary.TryGetValue("self", out var self) && self != null) + { + links.Self = self.ToString(); + } + + if (linksDictionary.TryGetValue("next", out var next) && next != null) + { + links.Next = next.ToString(); + } + + if (linksDictionary.TryGetValue("prev", out var prev) && prev != null) + { + links.Prev = prev.ToString(); + } + + result.Links = links; + } + } + + return result; + } +} \ No newline at end of file diff --git a/src/Api/PubnubApi/JsonDataParse/PNDataSyncEntityResultJsonDataParse.cs b/src/Api/PubnubApi/JsonDataParse/PNDataSyncEntityResultJsonDataParse.cs new file mode 100644 index 000000000..b70f40c9e --- /dev/null +++ b/src/Api/PubnubApi/JsonDataParse/PNDataSyncEntityResultJsonDataParse.cs @@ -0,0 +1,59 @@ +using System; +using System.Collections.Generic; +using PubnubApi.EndPoint; + +namespace PubnubApi; + +internal static class PNDataSyncEntityResultJsonDataParse +{ + internal static PNDataSyncEntityResult GetObject(IJsonPluggableLibrary jsonPlug, List listObject) + { + Dictionary dictionaryObject = (listObject != null && listObject.Count == 1) + ? jsonPlug.ConvertToDictionaryObject(listObject[0]) + : null; + PNDataSyncEntityResult result = null; + if (dictionaryObject != null && dictionaryObject.ContainsKey("data")) + { + result = new PNDataSyncEntityResult(); + + Dictionary dataDictionary = + jsonPlug.ConvertToDictionaryObject(dictionaryObject["data"]); + if (dataDictionary != null && dataDictionary.Count > 0) + { + result.Id = dataDictionary.ContainsKey("id") && dataDictionary["id"] != null + ? dataDictionary["id"].ToString() + : null; + result.EntityClass = dataDictionary.ContainsKey("entityClass") && dataDictionary["entityClass"] != null + ? dataDictionary["entityClass"].ToString() + : null; + result.EntityClassVersion = dataDictionary.ContainsKey("entityClassVersion") && dataDictionary["entityClassVersion"] != null + ? (int)(long)dataDictionary["entityClassVersion"] + : 1; + result.Status = dataDictionary.ContainsKey("status") && dataDictionary["status"] != null + ? dataDictionary["status"].ToString() + : null; + result.CreatedAt = dataDictionary.ContainsKey("createdAt") && dataDictionary["createdAt"] != null + ? dataDictionary["createdAt"].ToString() + : null; + result.UpdatedAt = dataDictionary.ContainsKey("updatedAt") && dataDictionary["updatedAt"] != null + ? dataDictionary["updatedAt"].ToString() + : null; + result.ETag = dataDictionary.ContainsKey("eTag") && dataDictionary["eTag"] != null + ? dataDictionary["eTag"].ToString() + : null; + result.ExpiresAt = dataDictionary.ContainsKey("expiresAt") && dataDictionary["expiresAt"] != null + ? dataDictionary["expiresAt"].ToString() + : null; + result.Id = dataDictionary.ContainsKey("id") && dataDictionary["id"] != null + ? dataDictionary["id"].ToString() + : null; + if (dataDictionary.TryGetValue("payload", out var payloadObject) && payloadObject != null) + { + result.Payload = jsonPlug.ConvertToDictionaryObject(payloadObject); + } + } + } + + return result; + } +} \ No newline at end of file diff --git a/src/Api/PubnubApi/Model/Consumer/PNStatus.cs b/src/Api/PubnubApi/Model/Consumer/PNStatus.cs index 361ece757..907af0250 100644 --- a/src/Api/PubnubApi/Model/Consumer/PNStatus.cs +++ b/src/Api/PubnubApi/Model/Consumer/PNStatus.cs @@ -521,6 +521,66 @@ public void Retry() } } break; + case PNOperationType.PNDataSyncCreateEntity: + if (savedEndpointOperation is CreateEntityOperation) + { + CreateEntityOperation endpoint = savedEndpointOperation as CreateEntityOperation; + if (endpoint != null) + { + endpoint.Retry(); + } + } + break; + case PNOperationType.PNDataSyncGetEntity: + if (savedEndpointOperation is GetEntityOperation) + { + GetEntityOperation endpoint = savedEndpointOperation as GetEntityOperation; + if (endpoint != null) + { + endpoint.Retry(); + } + } + break; + case PNOperationType.PNDataSyncGetEntities: + if (savedEndpointOperation is GetEntitiesOperation) + { + GetEntitiesOperation endpoint = savedEndpointOperation as GetEntitiesOperation; + if (endpoint != null) + { + endpoint.Retry(); + } + } + break; + case PNOperationType.PNDataSyncUpdateEntity: + if (savedEndpointOperation is UpdateEntityOperation) + { + UpdateEntityOperation endpoint = savedEndpointOperation as UpdateEntityOperation; + if (endpoint != null) + { + endpoint.Retry(); + } + } + break; + case PNOperationType.PNDataSyncPatchEntity: + if (savedEndpointOperation is PatchEntityOperation) + { + PatchEntityOperation endpoint = savedEndpointOperation as PatchEntityOperation; + if (endpoint != null) + { + endpoint.Retry(); + } + } + break; + case PNOperationType.PNDataSyncDeleteEntity: + if (savedEndpointOperation is DeleteEntityOperation) + { + DeleteEntityOperation endpoint = savedEndpointOperation as DeleteEntityOperation; + if (endpoint != null) + { + endpoint.Retry(); + } + } + break; default: break; } diff --git a/src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncDeleteEntityResultExt.cs b/src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncDeleteEntityResultExt.cs new file mode 100644 index 000000000..024f9ee2d --- /dev/null +++ b/src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncDeleteEntityResultExt.cs @@ -0,0 +1,19 @@ +using System; +using PubnubApi.EndPoint; + +namespace PubnubApi; + +public class PNDataSyncDeleteEntityResultExt : PNCallback +{ + readonly Action callbackAction; + + public PNDataSyncDeleteEntityResultExt(Action callback) + { + this.callbackAction = callback; + } + + public override void OnResponse(PNDataSyncDeleteEntityResult result, PNStatus status) + { + callbackAction?.Invoke(result, status); + } +} \ No newline at end of file diff --git a/src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncEntitiesListResultExt.cs b/src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncEntitiesListResultExt.cs new file mode 100644 index 000000000..c4938c3c8 --- /dev/null +++ b/src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncEntitiesListResultExt.cs @@ -0,0 +1,19 @@ +using System; +using PubnubApi.EndPoint; + +namespace PubnubApi; + +public class PNDataSyncEntitiesListResultExt : PNCallback +{ + readonly Action callbackAction; + + public PNDataSyncEntitiesListResultExt(Action callback) + { + this.callbackAction = callback; + } + + public override void OnResponse(PNDataSyncEntitiesListResult result, PNStatus status) + { + callbackAction?.Invoke(result, status); + } +} \ No newline at end of file diff --git a/src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncEntityResultExt.cs b/src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncEntityResultExt.cs new file mode 100644 index 000000000..76363d3a4 --- /dev/null +++ b/src/Api/PubnubApi/Model/Derived/DataSync/PNDataSyncEntityResultExt.cs @@ -0,0 +1,19 @@ +using System; +using PubnubApi.EndPoint; + +namespace PubnubApi; + +public class PNDataSyncEntityResultExt : PNCallback +{ + readonly Action callbackAction; + + public PNDataSyncEntityResultExt(Action callback) + { + this.callbackAction = callback; + } + + public override void OnResponse(PNDataSyncEntityResult result, PNStatus status) + { + callbackAction?.Invoke(result, status); + } +} \ No newline at end of file diff --git a/src/Api/PubnubApi/Pubnub.cs b/src/Api/PubnubApi/Pubnub.cs index 2fee9d258..95f5cd3f6 100644 --- a/src/Api/PubnubApi/Pubnub.cs +++ b/src/Api/PubnubApi/Pubnub.cs @@ -1185,6 +1185,8 @@ public void SetLogger(IPubnubLogger logger) #region "Properties" + public DataSync DataSync { get; } + public IPubnubUnitTest PubnubUnitTest { get => pubnubUnitTest; @@ -1265,6 +1267,7 @@ public Pubnub(PNConfiguration config, IHttpClientService httpTransportService = httpTransportService ?? new HttpClientService(proxy: config.Proxy); httpClientService.SetLogger(logger); transportMiddleware = middleware ?? new Middleware(httpClientService, config, this, tokenManager); + DataSync = new DataSync(this, pubnubUnitTest, tokenManager); logger?.Debug(GetConfigurationLogString(config)); } diff --git a/src/Api/PubnubApi/PubnubCoreBase.cs b/src/Api/PubnubApi/PubnubCoreBase.cs index 7d14a3580..2ee9ba37b 100644 --- a/src/Api/PubnubApi/PubnubCoreBase.cs +++ b/src/Api/PubnubApi/PubnubCoreBase.cs @@ -1055,6 +1055,20 @@ protected PNStatus GetStatusIfError(RequestState asyncRequestState, string } } } + else if (deserializeStatus.TryGetValue("errors", out var errorListObject)) + { + var aggregateErrorsString = errorListObject.ToString(); + if (pubnubConfig.TryGetValue(PubnubInstance.InstanceId, out currentConfig)) + { + statusCode = asyncRequestState?.Response?.StatusCode ?? 500; + status = new StatusBuilder(currentConfig, jsonLib).CreateStatusResponse( + type, + PNStatusCategory.PNUnknownCategory, + asyncRequestState, + statusCode, + new PNException(aggregateErrorsString)); + } + } } else if (jsonString.ToLowerInvariant().TrimStart().IndexOf(" PostRequest(TransportRequest transportReque HttpContent postData = null; if (!string.IsNullOrEmpty(transportRequest.BodyContentString)) { - postData = new StringContent(transportRequest.BodyContentString, Encoding.UTF8, "application/json"); + var contentType = "application/json"; + if (transportRequest.Headers.TryGetValue("Content-Type", out var ct)) + { + contentType = ct; + } + postData = new StringContent(transportRequest.BodyContentString, Encoding.UTF8); + postData.Headers.ContentType = System.Net.Http.Headers.MediaTypeHeaderValue.Parse(contentType); } else if (transportRequest.BodyContentBytes != null) { @@ -119,6 +125,13 @@ public async Task PostRequest(TransportRequest transportReque HttpRequestMessage requestMessage = new HttpRequestMessage(method: HttpMethod.Post, requestUri: transportRequest.RequestUrl) { Content = postData }; + foreach (var kvp in transportRequest.Headers) + { + if (!string.Equals(kvp.Key, "Content-Type", StringComparison.OrdinalIgnoreCase)) + { + requestMessage.Headers.Add(kvp.Key, kvp.Value); + } + } logger?.Debug( $"HttpClient Service:Sending http request {transportRequest.RequestType} to {transportRequest.RequestUrl}" + (requestMessage.Headers.Any() @@ -176,7 +189,13 @@ public async Task PutRequest(TransportRequest transportReques if (!string.IsNullOrEmpty(transportRequest.BodyContentString)) { - putData = new StringContent(transportRequest.BodyContentString, Encoding.UTF8, "application/json"); + var contentType = "application/json"; + if (transportRequest.Headers.TryGetValue("Content-Type", out var ct)) + { + contentType = ct; + } + putData = new StringContent(transportRequest.BodyContentString, Encoding.UTF8); + putData.Headers.ContentType = System.Net.Http.Headers.MediaTypeHeaderValue.Parse(contentType); } else if (transportRequest.BodyContentBytes != null) { @@ -190,9 +209,9 @@ public async Task PutRequest(TransportRequest transportReques HttpRequestMessage requestMessage = new HttpRequestMessage(method: HttpMethod.Put, requestUri: transportRequest.RequestUrl) { Content = putData }; - if (transportRequest.Headers.Keys.Count > 0) + foreach (var kvp in transportRequest.Headers) { - foreach (var kvp in transportRequest.Headers) + if (!string.Equals(kvp.Key, "Content-Type", StringComparison.OrdinalIgnoreCase)) { requestMessage.Headers.Add(kvp.Key, kvp.Value); } @@ -318,8 +337,13 @@ public async Task PatchRequest(TransportRequest transportRequ if (!string.IsNullOrEmpty(transportRequest.BodyContentString)) { - patchData = new StringContent(transportRequest.BodyContentString, Encoding.UTF8, - "application/json"); + var contentType = "application/json"; + if (transportRequest.Headers.TryGetValue("Content-Type", out var ct)) + { + contentType = ct; + } + patchData = new StringContent(transportRequest.BodyContentString, Encoding.UTF8); + patchData.Headers.ContentType = System.Net.Http.Headers.MediaTypeHeaderValue.Parse(contentType); } else if (transportRequest.BodyContentBytes != null) { @@ -333,11 +357,11 @@ public async Task PatchRequest(TransportRequest transportRequ HttpRequestMessage requestMessage = new HttpRequestMessage(new HttpMethod("PATCH"), requestUri: transportRequest.RequestUrl) { Content = patchData }; - if (transportRequest.Headers.Keys.Count > 0) + foreach (var kvp in transportRequest.Headers) { - foreach (var kvp in transportRequest.Headers) + if (!string.Equals(kvp.Key, "Content-Type", StringComparison.OrdinalIgnoreCase)) { - requestMessage.Headers.Add(kvp.Key, $"\"{kvp.Value}\""); + requestMessage.Headers.Add(kvp.Key, kvp.Value); } } diff --git a/src/Api/PubnubApiPCL/PubnubApiPCL.csproj b/src/Api/PubnubApiPCL/PubnubApiPCL.csproj index 2f23d4608..ea5d78455 100644 --- a/src/Api/PubnubApiPCL/PubnubApiPCL.csproj +++ b/src/Api/PubnubApiPCL/PubnubApiPCL.csproj @@ -79,6 +79,51 @@ EndPoint\ChannelGroup\RemoveChannelsFromChannelGroupOperation.cs + + EndPoint\DataSync\CreateEntityOperation.cs + + + EndPoint\DataSync\CreateEntityParameters.cs + + + EndPoint\DataSync\DataSync.cs + + + EndPoint\DataSync\DeleteEntityOperation.cs + + + EndPoint\DataSync\DeleteEntityParameters.cs + + + EndPoint\DataSync\GetEntitiesOperation.cs + + + EndPoint\DataSync\GetEntitiesParameters.cs + + + EndPoint\DataSync\GetEntityOperation.cs + + + EndPoint\DataSync\GetEntityParameters.cs + + + EndPoint\DataSync\PatchEntityOperation.cs + + + EndPoint\DataSync\PatchEntityParameters.cs + + + EndPoint\DataSync\PNDataSyncDeleteEntityResult.cs + + + EndPoint\DataSync\PNDataSyncEntityResult.cs + + + EndPoint\DataSync\UpdateEntityOperation.cs + + + EndPoint\DataSync\UpdateEntityParameters.cs + @@ -291,6 +336,12 @@ + + JsonDataParse\PNDataSyncEntitiesListResultJsonDataParse.cs + + + JsonDataParse\PNDataSyncEntityResultJsonDataParse.cs + @@ -472,6 +523,15 @@ Model\Derived\ChannelGroup\PNChannelGroupsRemoveChannelResultExt.cs + + Model\Derived\DataSync\PNDataSyncDeleteEntityResultExt.cs + + + Model\Derived\DataSync\PNDataSyncEntitiesListResultExt.cs + + + Model\Derived\DataSync\PNDataSyncEntityResultExt.cs + @@ -695,5 +755,29 @@ + + + + EndPoint\DataSync\C# - DXD data sync Entity API.md + + + EndPoint\DataSync\C# - DXD data sync Entity Class API.md + + + EndPoint\DataSync\C# - DXD data sync User API.md + + + EndPoint\DataSync\common.yaml + + + EndPoint\DataSync\test-entity-api.ps1 + + + EndPoint\DataSync\v4-data.yaml + + + EndPoint\DataSync\v4-metadata.yaml + + diff --git a/src/Api/PubnubApiUWP/PubnubApiUWP.csproj b/src/Api/PubnubApiUWP/PubnubApiUWP.csproj index 113a10ce3..0f45a4c38 100644 --- a/src/Api/PubnubApiUWP/PubnubApiUWP.csproj +++ b/src/Api/PubnubApiUWP/PubnubApiUWP.csproj @@ -204,6 +204,51 @@ EndPoint\ChannelGroup\RemoveChannelsFromChannelGroupOperation.cs + + EndPoint\DataSync\CreateEntityOperation.cs + + + EndPoint\DataSync\CreateEntityParameters.cs + + + EndPoint\DataSync\DataSync.cs + + + EndPoint\DataSync\DeleteEntityOperation.cs + + + EndPoint\DataSync\DeleteEntityParameters.cs + + + EndPoint\DataSync\GetEntitiesOperation.cs + + + EndPoint\DataSync\GetEntitiesParameters.cs + + + EndPoint\DataSync\GetEntityOperation.cs + + + EndPoint\DataSync\GetEntityParameters.cs + + + EndPoint\DataSync\PatchEntityOperation.cs + + + EndPoint\DataSync\PatchEntityParameters.cs + + + EndPoint\DataSync\PNDataSyncDeleteEntityResult.cs + + + EndPoint\DataSync\PNDataSyncEntityResult.cs + + + EndPoint\DataSync\UpdateEntityOperation.cs + + + EndPoint\DataSync\UpdateEntityParameters.cs + @@ -420,6 +465,12 @@ + + JsonDataParse\PNDataSyncEntitiesListResultJsonDataParse.cs + + + JsonDataParse\PNDataSyncEntityResultJsonDataParse.cs + @@ -601,6 +652,15 @@ Model\Derived\ChannelGroup\PNChannelGroupsRemoveChannelResultExt.cs + + Model\Derived\DataSync\PNDataSyncDeleteEntityResultExt.cs + + + Model\Derived\DataSync\PNDataSyncEntitiesListResultExt.cs + + + Model\Derived\DataSync\PNDataSyncEntityResultExt.cs + @@ -813,6 +873,29 @@ + + + EndPoint\DataSync\C# - DXD data sync Entity API.md + + + EndPoint\DataSync\C# - DXD data sync Entity Class API.md + + + EndPoint\DataSync\C# - DXD data sync User API.md + + + EndPoint\DataSync\common.yaml + + + EndPoint\DataSync\test-entity-api.ps1 + + + EndPoint\DataSync\v4-data.yaml + + + EndPoint\DataSync\v4-metadata.yaml + + 14.0 diff --git a/src/Api/PubnubApiUnity/PubnubApiUnity.csproj b/src/Api/PubnubApiUnity/PubnubApiUnity.csproj index a804de806..6b887be9e 100644 --- a/src/Api/PubnubApiUnity/PubnubApiUnity.csproj +++ b/src/Api/PubnubApiUnity/PubnubApiUnity.csproj @@ -91,6 +91,51 @@ EndPoint\ChannelGroup\RemoveChannelsFromChannelGroupOperation.cs + + EndPoint\DataSync\CreateEntityOperation.cs + + + EndPoint\DataSync\CreateEntityParameters.cs + + + EndPoint\DataSync\DataSync.cs + + + EndPoint\DataSync\DeleteEntityOperation.cs + + + EndPoint\DataSync\DeleteEntityParameters.cs + + + EndPoint\DataSync\GetEntitiesOperation.cs + + + EndPoint\DataSync\GetEntitiesParameters.cs + + + EndPoint\DataSync\GetEntityOperation.cs + + + EndPoint\DataSync\GetEntityParameters.cs + + + EndPoint\DataSync\PatchEntityOperation.cs + + + EndPoint\DataSync\PatchEntityParameters.cs + + + EndPoint\DataSync\PNDataSyncDeleteEntityResult.cs + + + EndPoint\DataSync\PNDataSyncEntityResult.cs + + + EndPoint\DataSync\UpdateEntityOperation.cs + + + EndPoint\DataSync\UpdateEntityParameters.cs + @@ -226,9 +271,24 @@ JsonDataParse\DeserializeToInternalObjectUtility.cs + + JsonDataParse\PNDataSyncEntitiesListResultJsonDataParse.cs + + + JsonDataParse\PNDataSyncEntityResultJsonDataParse.cs + Model\Consumer\Objects\PNMembershipMetadataResult.cs + + Model\Derived\DataSync\PNDataSyncDeleteEntityResultExt.cs + + + Model\Derived\DataSync\PNDataSyncEntitiesListResultExt.cs + + + Model\Derived\DataSync\PNDataSyncEntityResultExt.cs + PNSDK\DotNetPNSDKSource.cs @@ -725,6 +785,30 @@ + + + EndPoint\DataSync\C# - DXD data sync Entity API.md + + + EndPoint\DataSync\C# - DXD data sync Entity Class API.md + + + EndPoint\DataSync\C# - DXD data sync User API.md + + + EndPoint\DataSync\common.yaml + + + EndPoint\DataSync\test-entity-api.ps1 + + + EndPoint\DataSync\v4-data.yaml + + + EndPoint\DataSync\v4-metadata.yaml + + + PubNub is a Massively Scalable Web Push Service for Web and Mobile Games. This is a cloud-based service for broadcasting messages to thousands of web and mobile clients simultaneously diff --git a/src/Api/PubnubApiPCL/PubnubApiPCL.csproj b/src/Api/PubnubApiPCL/PubnubApiPCL.csproj index 49f3fcf55..d62b53cd8 100644 --- a/src/Api/PubnubApiPCL/PubnubApiPCL.csproj +++ b/src/Api/PubnubApiPCL/PubnubApiPCL.csproj @@ -14,7 +14,7 @@ PubnubPCL - 8.3.4 + 9.0.0 PubNub C# .NET - Web Data Push API Pandu Masabathula PubNub @@ -22,9 +22,7 @@ http://pubnub.s3.amazonaws.com/2011/powered-by-pubnub/pubnub-icon-600x600.png true https://github.com/pubnub/c-sharp/ - Fixed issue related to Reconnect to retain subscription state. -Removed value cap from maxRetry in Linear/Exponential policies. -Emit status PNUnexpectedDisconnectCategory when retry attempts exhausted. Instead every retry attempt failures. + Added Data Sync feature support. Web Data Push Real-time Notifications ESB Message Broadcasting Distributed Computing PubNub is a Massively Scalable Web Push Service for Web and Mobile Games. This is a cloud-based service for broadcasting messages to thousands of web and mobile clients simultaneously diff --git a/src/Api/PubnubApiUWP/PubnubApiUWP.csproj b/src/Api/PubnubApiUWP/PubnubApiUWP.csproj index 97a12ee09..207f9bec5 100644 --- a/src/Api/PubnubApiUWP/PubnubApiUWP.csproj +++ b/src/Api/PubnubApiUWP/PubnubApiUWP.csproj @@ -16,7 +16,7 @@ PubnubUWP - 8.3.4 + 9.0.0 PubNub C# .NET - Web Data Push API Pandu Masabathula PubNub @@ -24,9 +24,7 @@ http://pubnub.s3.amazonaws.com/2011/powered-by-pubnub/pubnub-icon-600x600.png true https://github.com/pubnub/c-sharp/ - Fixed issue related to Reconnect to retain subscription state. -Removed value cap from maxRetry in Linear/Exponential policies. -Emit status PNUnexpectedDisconnectCategory when retry attempts exhausted. Instead every retry attempt failures. + Added Data Sync feature support. Web Data Push Real-time Notifications ESB Message Broadcasting Distributed Computing PubNub is a Massively Scalable Web Push Service for Web and Mobile Games. This is a cloud-based service for broadcasting messages to thousands of web and mobile clients simultaneously diff --git a/src/Api/PubnubApiUnity/PubnubApiUnity.csproj b/src/Api/PubnubApiUnity/PubnubApiUnity.csproj index 822437344..17ccf6340 100644 --- a/src/Api/PubnubApiUnity/PubnubApiUnity.csproj +++ b/src/Api/PubnubApiUnity/PubnubApiUnity.csproj @@ -15,7 +15,7 @@ PubnubApiUnity - 8.3.4 + 9.0.0 PubNub C# .NET - Web Data Push API Pandu Masabathula PubNub From afda7fc61632065753693b7ccc7c6d9c7eecd373 Mon Sep 17 00:00:00 2001 From: "PUBNUB\\jakub.grzesiowski" Date: Tue, 8 Sep 2026 14:13:37 +0200 Subject: [PATCH 27/31] Fix compilation errors in tests --- .../WhenDataSyncMembershipEventIsReceived.cs | 4 ++-- .../DataSync/WhenDataSyncMembershipIsRequested.cs | 12 ++++++------ 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipEventIsReceived.cs b/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipEventIsReceived.cs index b5ca0eddd..219a24185 100644 --- a/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipEventIsReceived.cs +++ b/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipEventIsReceived.cs @@ -133,7 +133,7 @@ private async Task CreateTestMembership( Id = id, ChannelId = channel.Id, UserId = user.Id, - RelationshipClassVersion = TestRelationshipClassVersion, + MembershipClassVersion = TestRelationshipClassVersion, Status = status, Payload = payload ?? new Dictionary { @@ -213,7 +213,7 @@ public async Task ThenCreatingMembershipShouldDeliverCreateEvent() Id = membershipId, ChannelId = channel.Id, UserId = user.Id, - RelationshipClassVersion = TestRelationshipClassVersion, + MembershipClassVersion = TestRelationshipClassVersion, Status = "active", Payload = new Dictionary { diff --git a/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipIsRequested.cs b/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipIsRequested.cs index 9901e95b0..3abf6caa4 100644 --- a/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipIsRequested.cs +++ b/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipIsRequested.cs @@ -140,7 +140,7 @@ private async Task CreateTestMembership( { ChannelId = channelId, UserId = userId, - RelationshipClassVersion = TestRelationshipClassVersion, + MembershipClassVersion = TestRelationshipClassVersion, Status = status, Payload = payload ?? new Dictionary { @@ -181,7 +181,7 @@ public async Task ThenCreateWithAllFieldsShouldReturnCreatedMembership() { ChannelId = channel.Id, UserId = user.Id, - RelationshipClassVersion = TestRelationshipClassVersion, + MembershipClassVersion = TestRelationshipClassVersion, Status = "active", Payload = payload, @@ -214,7 +214,7 @@ public async Task ThenCreateWithExplicitIdShouldUseProvidedId() Id = membershipId, ChannelId = channel.Id, UserId = user.Id, - RelationshipClassVersion = TestRelationshipClassVersion, + MembershipClassVersion = TestRelationshipClassVersion, Status = "active", Payload = new Dictionary { { "key", "value" } }, @@ -237,7 +237,7 @@ public async Task ThenCreateWithoutIdShouldReturnServerGeneratedId() { ChannelId = channel.Id, UserId = user.Id, - RelationshipClassVersion = TestRelationshipClassVersion, + MembershipClassVersion = TestRelationshipClassVersion, Payload = new Dictionary { { "key", "value" } }, }); @@ -260,7 +260,7 @@ public async Task ThenCreateWithMinimalFieldsShouldSucceed() { ChannelId = channel.Id, UserId = user.Id, - RelationshipClassVersion = TestRelationshipClassVersion, + MembershipClassVersion = TestRelationshipClassVersion, }); @@ -994,7 +994,7 @@ public async Task ThenFullCrudLifecycleShouldSucceed() { ChannelId = channel.Id, UserId = user.Id, - RelationshipClassVersion = TestRelationshipClassVersion, + MembershipClassVersion = TestRelationshipClassVersion, Status = "new", Payload = new Dictionary { From 5bb3ae8c1aa97588ae1d9f237831b3414a008c0c Mon Sep 17 00:00:00 2001 From: "PUBNUB\\jakub.grzesiowski" Date: Tue, 8 Sep 2026 15:44:22 +0200 Subject: [PATCH 28/31] Tweak test delay --- src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs b/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs index d35cc8d9d..8a3da7aea 100644 --- a/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs +++ b/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs @@ -843,7 +843,7 @@ public static void IfHereNowIsCalledThenItShouldReturnInfoCipher() pubnub.Unsubscribe().Channels(new[] { channel }).Execute(); - if (!PubnubCommon.EnableStubTest) Thread.Sleep(1000); + if (!PubnubCommon.EnableStubTest) Thread.Sleep(3000); else Thread.Sleep(100); } From 134cd773d4d93ea26b3f8d767cfe6cbec671c2fd Mon Sep 17 00:00:00 2001 From: "PUBNUB\\jakub.grzesiowski" Date: Tue, 8 Sep 2026 15:56:14 +0200 Subject: [PATCH 29/31] More minor test tweaks --- src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs b/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs index 8a3da7aea..ae1dc57c1 100644 --- a/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs +++ b/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs @@ -578,7 +578,7 @@ public static async Task IfWithAsyncHereNowIsCalledThenItShouldReturnInfo() if (!receivedErrorMessage) { - if (!PubnubCommon.EnableStubTest) { Thread.Sleep(2000); } + if (!PubnubCommon.EnableStubTest) { Thread.Sleep(4000); } else Thread.Sleep(200); expected = "{\"status\": 200, \"message\": \"OK\", \"service\": \"Presence\", \"uuids\": [\"mytestuuid\"], \"occupancy\": 1}"; @@ -610,7 +610,7 @@ public static async Task IfWithAsyncHereNowIsCalledThenItShouldReturnInfo() pubnub.Unsubscribe().Channels(new[] { channel }).Execute(); - if (!PubnubCommon.EnableStubTest) Thread.Sleep(1000); + if (!PubnubCommon.EnableStubTest) Thread.Sleep(4000); else Thread.Sleep(100); } From 8776c8962a5f99d5c264cc5c869a83d0c61adefb Mon Sep 17 00:00:00 2001 From: "PUBNUB\\jakub.grzesiowski" Date: Tue, 8 Sep 2026 16:12:36 +0200 Subject: [PATCH 30/31] Another tweak --- src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs b/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs index ae1dc57c1..6d960af6a 100644 --- a/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs +++ b/src/UnitTests/PubnubApi.Tests/WhenAClientIsPresented.cs @@ -2614,7 +2614,7 @@ public static async Task ThenWhereNowShouldReturnSubscribedChannel() pubnub = createPubNubInstance(config, authToken); pubnub.Subscribe().Channels(new[] { subscribedChannel }).WithPresence().Execute(); - await Task.Delay(3000); + await Task.Delay(6000); ManualResetEvent whereNowManualEvent = new ManualResetEvent(false); PNWhereNowResult whereNowResult = null; From cabddcbe3f1c1e57b7c8fc690385f666311ec8c1 Mon Sep 17 00:00:00 2001 From: "PUBNUB\\jakub.grzesiowski" Date: Tue, 8 Sep 2026 17:48:31 +0200 Subject: [PATCH 31/31] Update event parsing for new channel/user/membership event types --- .../JsonDataParse/PNDataSyncEventJsonDataParse.cs | 8 ++++---- .../Model/Consumer/Pubsub/PNDataSyncEventResult.cs | 4 ++-- .../DataSync/WhenDataSyncChannelEventIsReceived.cs | 8 ++++---- .../DataSync/WhenDataSyncMembershipEventIsReceived.cs | 8 ++++---- .../DataSync/WhenDataSyncUserEventIsReceived.cs | 8 ++++---- 5 files changed, 18 insertions(+), 18 deletions(-) diff --git a/src/Api/PubnubApi/JsonDataParse/PNDataSyncEventJsonDataParse.cs b/src/Api/PubnubApi/JsonDataParse/PNDataSyncEventJsonDataParse.cs index f77cac8f9..8ebf325af 100644 --- a/src/Api/PubnubApi/JsonDataParse/PNDataSyncEventJsonDataParse.cs +++ b/src/Api/PubnubApi/JsonDataParse/PNDataSyncEventJsonDataParse.cs @@ -50,7 +50,7 @@ internal static PNDataSyncEventResult GetObject(IJsonPluggableLibrary jsonPlug, result.Id = GetStringValue(data, "id"); result.DeletedAt = GetStringValue(data, "deletedAt"); } - else if (type == "entity") + else if (type == "entity" || type == "user" || type == "channel") { result.EntityData = new PNDataSyncEntityResult { @@ -67,13 +67,13 @@ internal static PNDataSyncEventResult GetObject(IJsonPluggableLibrary jsonPlug, ExpiresAt = GetStringValue(data, "expiresAt"), }; } - else if (type == "relationship") + else if (type == "relationship" || type == "membership") { result.RelationshipData = new PNDataSyncRelationshipResult { Id = GetStringValue(data, "id"), - EntityAId = GetStringValue(data, "entityAId"), - EntityBId = GetStringValue(data, "entityBId"), + EntityAId = GetStringValue(data, type == "membership" ? "channelId" : "entityAId"), + EntityBId = GetStringValue(data, type == "membership" ? "userId" : "entityBId"), RelationshipClass = result.ClassName, RelationshipClassVersion = result.ClassVersion, Status = GetStringValue(data, "status"), diff --git a/src/Api/PubnubApi/Model/Consumer/Pubsub/PNDataSyncEventResult.cs b/src/Api/PubnubApi/Model/Consumer/Pubsub/PNDataSyncEventResult.cs index 5bf50a503..df40c6125 100644 --- a/src/Api/PubnubApi/Model/Consumer/Pubsub/PNDataSyncEventResult.cs +++ b/src/Api/PubnubApi/Model/Consumer/Pubsub/PNDataSyncEventResult.cs @@ -9,8 +9,8 @@ public class PNDataSyncEventResult public string ClassName { get; internal set; } = ""; public int ClassVersion { get; internal set; } public string ClassLevel { get; internal set; } = ""; //values = Global/SubKey; distinguishes built-in classes from developer-defined ones of the same name - public PubnubApi.EndPoint.PNDataSyncEntityResult EntityData { get; internal set; } //Populated when Type = entity (create/update) - public PubnubApi.EndPoint.PNDataSyncRelationshipResult RelationshipData { get; internal set; } //Populated when Type = relationship (create/update) + public PubnubApi.EndPoint.PNDataSyncEntityResult EntityData { get; internal set; } //Populated when Type = entity/user/channel (create/update) + public PubnubApi.EndPoint.PNDataSyncRelationshipResult RelationshipData { get; internal set; } //Populated when Type = relationship/membership (create/update) public string Id { get; internal set; } //Populated for delete events public string DeletedAt { get; internal set; } //Populated for delete events public long Timestamp { get; internal set; } diff --git a/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncChannelEventIsReceived.cs b/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncChannelEventIsReceived.cs index 7fa8731e6..d59f6576d 100644 --- a/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncChannelEventIsReceived.cs +++ b/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncChannelEventIsReceived.cs @@ -147,7 +147,7 @@ public async Task ThenCreatingChannelShouldDeliverCreateEvent() Assert.That(dataSyncEvent.Event, Is.EqualTo("create").IgnoreCase); Assert.That(dataSyncEvent.Source, Is.EqualTo("data-sync")); - Assert.That(dataSyncEvent.Type, Is.EqualTo("entity").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("channel").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(channelId)); Assert.That(dataSyncEvent.EntityData, Is.Not.Null); Assert.That(dataSyncEvent.EntityData.Id, Is.EqualTo(channelId)); @@ -178,7 +178,7 @@ public async Task ThenUpdatingChannelShouldDeliverUpdateEvent() Assert.That(dataSyncEvent.Event, Is.EqualTo("update").IgnoreCase); Assert.That(dataSyncEvent.Source, Is.EqualTo("data-sync")); - Assert.That(dataSyncEvent.Type, Is.EqualTo("entity").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("channel").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(created.Id)); Assert.That(dataSyncEvent.EntityData, Is.Not.Null); Assert.That(dataSyncEvent.EntityData.Id, Is.EqualTo(created.Id)); @@ -213,7 +213,7 @@ public async Task ThenPatchingChannelShouldDeliverUpdateEvent() }); Assert.That(dataSyncEvent.Event, Is.EqualTo("update").IgnoreCase); - Assert.That(dataSyncEvent.Type, Is.EqualTo("entity").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("channel").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(created.Id)); Assert.That(dataSyncEvent.EntityData, Is.Not.Null); Assert.That(dataSyncEvent.EntityData.Id, Is.EqualTo(created.Id)); @@ -241,7 +241,7 @@ public async Task ThenDeletingChannelShouldDeliverDeleteEvent() Assert.That(dataSyncEvent.Event, Is.EqualTo("delete").IgnoreCase); Assert.That(dataSyncEvent.Source, Is.EqualTo("data-sync")); - Assert.That(dataSyncEvent.Type, Is.EqualTo("entity").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("channel").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(created.Id)); Assert.That(dataSyncEvent.Id, Is.EqualTo(created.Id)); Assert.That(dataSyncEvent.DeletedAt, Is.Not.Null.And.Not.Empty); diff --git a/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipEventIsReceived.cs b/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipEventIsReceived.cs index 219a24185..47bce5cb2 100644 --- a/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipEventIsReceived.cs +++ b/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncMembershipEventIsReceived.cs @@ -229,7 +229,7 @@ public async Task ThenCreatingMembershipShouldDeliverCreateEvent() Assert.That(dataSyncEvent.Event, Is.EqualTo("create").IgnoreCase); Assert.That(dataSyncEvent.Source, Is.EqualTo("data-sync")); - Assert.That(dataSyncEvent.Type, Is.EqualTo("relationship").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("membership").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(user.Id)); Assert.That(dataSyncEvent.RelationshipData, Is.Not.Null); Assert.That(dataSyncEvent.RelationshipData.Id, Is.EqualTo(membershipId)); @@ -267,7 +267,7 @@ public async Task ThenUpdatingMembershipShouldDeliverUpdateEvent() Assert.That(dataSyncEvent.Event, Is.EqualTo("update").IgnoreCase); Assert.That(dataSyncEvent.Source, Is.EqualTo("data-sync")); - Assert.That(dataSyncEvent.Type, Is.EqualTo("relationship").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("membership").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(created.UserId)); Assert.That(dataSyncEvent.RelationshipData, Is.Not.Null); Assert.That(dataSyncEvent.RelationshipData.Id, Is.EqualTo(created.Id)); @@ -303,7 +303,7 @@ public async Task ThenPatchingMembershipShouldDeliverUpdateEvent() }); Assert.That(dataSyncEvent.Event, Is.EqualTo("update").IgnoreCase); - Assert.That(dataSyncEvent.Type, Is.EqualTo("relationship").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("membership").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(created.UserId)); Assert.That(dataSyncEvent.RelationshipData, Is.Not.Null); Assert.That(dataSyncEvent.RelationshipData.Id, Is.EqualTo(created.Id)); @@ -331,7 +331,7 @@ public async Task ThenDeletingMembershipShouldDeliverDeleteEvent() Assert.That(dataSyncEvent.Event, Is.EqualTo("delete").IgnoreCase); Assert.That(dataSyncEvent.Source, Is.EqualTo("data-sync")); - Assert.That(dataSyncEvent.Type, Is.EqualTo("relationship").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("membership").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(created.UserId)); Assert.That(dataSyncEvent.Id, Is.EqualTo(created.Id)); Assert.That(dataSyncEvent.DeletedAt, Is.Not.Null.And.Not.Empty); diff --git a/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncUserEventIsReceived.cs b/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncUserEventIsReceived.cs index 95575a19b..566e64772 100644 --- a/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncUserEventIsReceived.cs +++ b/src/UnitTests/PubnubApi.Tests/DataSync/WhenDataSyncUserEventIsReceived.cs @@ -146,7 +146,7 @@ public async Task ThenCreatingUserShouldDeliverCreateEvent() Assert.That(dataSyncEvent.Event, Is.EqualTo("create").IgnoreCase); Assert.That(dataSyncEvent.Source, Is.EqualTo("data-sync")); - Assert.That(dataSyncEvent.Type, Is.EqualTo("entity").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("user").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(userId)); Assert.That(dataSyncEvent.EntityData, Is.Not.Null); Assert.That(dataSyncEvent.EntityData.Id, Is.EqualTo(userId)); @@ -177,7 +177,7 @@ public async Task ThenUpdatingUserShouldDeliverUpdateEvent() Assert.That(dataSyncEvent.Event, Is.EqualTo("update").IgnoreCase); Assert.That(dataSyncEvent.Source, Is.EqualTo("data-sync")); - Assert.That(dataSyncEvent.Type, Is.EqualTo("entity").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("user").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(created.Id)); Assert.That(dataSyncEvent.EntityData, Is.Not.Null); Assert.That(dataSyncEvent.EntityData.Id, Is.EqualTo(created.Id)); @@ -213,7 +213,7 @@ public async Task ThenPatchingUserShouldDeliverUpdateEvent() }); Assert.That(dataSyncEvent.Event, Is.EqualTo("update").IgnoreCase); - Assert.That(dataSyncEvent.Type, Is.EqualTo("entity").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("user").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(created.Id)); Assert.That(dataSyncEvent.EntityData, Is.Not.Null); Assert.That(dataSyncEvent.EntityData.Id, Is.EqualTo(created.Id)); @@ -241,7 +241,7 @@ public async Task ThenDeletingUserShouldDeliverDeleteEvent() Assert.That(dataSyncEvent.Event, Is.EqualTo("delete").IgnoreCase); Assert.That(dataSyncEvent.Source, Is.EqualTo("data-sync")); - Assert.That(dataSyncEvent.Type, Is.EqualTo("entity").IgnoreCase); + Assert.That(dataSyncEvent.Type, Is.EqualTo("user").IgnoreCase); Assert.That(dataSyncEvent.Channel, Is.EqualTo(created.Id)); Assert.That(dataSyncEvent.Id, Is.EqualTo(created.Id)); Assert.That(dataSyncEvent.DeletedAt, Is.Not.Null.And.Not.Empty);