first commit
This commit is contained in:
@@ -0,0 +1,110 @@
|
||||
using MassTransit;
|
||||
using OpenWarehouse.warehouse.api.Common.ItemTagManager;
|
||||
using OpenWarehouse.warehouse.api.Common.Kafka;
|
||||
using OpenWarehouse.warehouse.api.Common.Kafka.TagTree;
|
||||
using OpenWarehouse.warehouse.api.Common.TaskManager;
|
||||
using OpenWarehouse.warehouse.api.Model.Item;
|
||||
using ItemCategoryTag = OpenWarehouse.warehouse.api.Model.Item.ItemCategoryTag;
|
||||
using TaskStatus = OpenWarehouse.warehouse.api.Common.TaskManager.TaskStatus;
|
||||
|
||||
namespace OpenWarehouse.warehouse.api.Services;
|
||||
|
||||
public class TagTreeConsumer(ILogger<TagTreeConsumer> logger, ITaskManager taskManager
|
||||
,IItemTagManager itemTagManager) : IConsumer<TagTreeKafkaMessage>
|
||||
{
|
||||
public async Task Consume(ConsumeContext<TagTreeKafkaMessage> context)
|
||||
{
|
||||
logger.LogInformation("Received and processing Kafka Message: {Guid}", context.Message.Guid);
|
||||
|
||||
var status = TaskStatus.Processing;
|
||||
|
||||
try
|
||||
{
|
||||
await taskManager.UpdateTaskStatus(context.Message.Guid, status);
|
||||
|
||||
switch (context.Message.Action)
|
||||
{
|
||||
case Actions.Create:
|
||||
status = await Create(context.Message);
|
||||
break;
|
||||
case Actions.Delete:
|
||||
status = await Delete(context.Message.ItemCategoryTagsList);
|
||||
break;
|
||||
default:
|
||||
logger.LogWarning("Action not allowed in request: {Guid}", context.Message.Guid);
|
||||
status = TaskStatus.Error;
|
||||
break;
|
||||
}
|
||||
|
||||
await taskManager.UpdateTaskStatus(context.Message.Guid, status);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
logger.LogError(ex, "Error occurred while processing Kafka message: {Guid}", context.Message.Guid);
|
||||
status = TaskStatus.Error;
|
||||
await taskManager.UpdateTaskStatus(context.Message.Guid, status);
|
||||
}
|
||||
}
|
||||
|
||||
private async Task<TaskStatus> Create(TagTreeKafkaMessage context)
|
||||
{
|
||||
try
|
||||
{
|
||||
var len = context.ItemCategoryTagsList.Count;
|
||||
var preProcessingData = context.ItemCategoryTagsList
|
||||
.Distinct()
|
||||
.ToList();
|
||||
var result = await itemTagManager.CreateTagTree(preProcessingData);
|
||||
|
||||
if (result == 0)
|
||||
{
|
||||
logger.LogWarning("Failed to create tag tree for message: {Guid}", context.Guid);
|
||||
return TaskStatus.Error;
|
||||
}
|
||||
else if (result < len)
|
||||
{
|
||||
logger.LogInformation("Partial success in creating tag tree for message: {Guid}", context.Guid);
|
||||
return TaskStatus.NotFullSuccess;
|
||||
}
|
||||
else
|
||||
{
|
||||
logger.LogInformation("Successfully created tag tree for message: {Guid}", context.Guid);
|
||||
return TaskStatus.Success;
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
logger.LogError(ex, "Error occurred while creating tag tree for message: {Guid}", context.Guid);
|
||||
return TaskStatus.Error;
|
||||
}
|
||||
}
|
||||
|
||||
private async Task<TaskStatus> Delete(List<ItemCategoryTag> names)
|
||||
{
|
||||
try
|
||||
{
|
||||
var result = await itemTagManager.DeleteTagsByNames(names);
|
||||
|
||||
if (result == 0)
|
||||
{
|
||||
logger.LogWarning("Failed to delete tags for names: {Names}", names);
|
||||
return TaskStatus.Error;
|
||||
}
|
||||
else if (result < names.Count)
|
||||
{
|
||||
logger.LogInformation("Partial success in deleting tags for names: {Names}", names);
|
||||
return TaskStatus.NotFullSuccess;
|
||||
}
|
||||
else
|
||||
{
|
||||
logger.LogInformation("Successfully deleted tags for names: {Names}", names);
|
||||
return TaskStatus.Success;
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
logger.LogError(ex, "Error occurred while deleting tags for names: {Names}", names);
|
||||
return TaskStatus.Error;
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user