111 lines
4.1 KiB
C#
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;
|
|
}
|
|
}
|
|
}
|