Files
OpenWarehouse/OpenWarehouse.warehouse.api/Services/TagTreeConsumer.cs
T
2025-02-26 23:42:19 +01:00

111 lines
4.1 KiB
C#

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;
}
}
}