OrderSubscriptionMessageGrpcService.java 6.56 KB
package com.tianting.infoloop.service.grpc;

import com.baomidou.dynamic.datasource.annotation.DS;
import com.infoloop.tianting.ordersubscriptionmessageservice.BatchDeleteOrderSubscriptionMessagesByIdsRpcRequest;
import com.infoloop.tianting.ordersubscriptionmessageservice.BatchDeleteOrderSubscriptionMessagesByIdsRpcResponse;
import com.infoloop.tianting.ordersubscriptionmessageservice.BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcRequest;
import com.infoloop.tianting.ordersubscriptionmessageservice.BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcResponse;
import com.infoloop.tianting.ordersubscriptionmessageservice.CreateOrderSubscriptionMessageRpcRequest;
import com.infoloop.tianting.ordersubscriptionmessageservice.CreateOrderSubscriptionMessageRpcResponse;
import com.infoloop.tianting.ordersubscriptionmessageservice.GetUserOrderSubscriptionMessageHistoryRpcRequest;
import com.infoloop.tianting.ordersubscriptionmessageservice.GetUserOrderSubscriptionMessageHistoryRpcResponse;
import com.infoloop.tianting.ordersubscriptionmessageservice.OrderSubscriptionMessageRpcResponse;
import com.infoloop.tianting.ordersubscriptionmessageservice.OrderSubscriptionMessageServiceRpcGrpc;
import com.tianting.infoloop.model.db.OrderSubscriptionMessageDb;
import com.tianting.infoloop.service.OrderSubscriptionMessageDbService;
import io.grpc.Status;
import io.grpc.stub.StreamObserver;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

import java.time.ZoneId;
import java.util.stream.Collectors;

import static com.tianting.infoloop.constants.ConfigConstants.DATA_SOURCE_MASTER;

@Slf4j
@Service
@RequiredArgsConstructor(onConstructor_ = @Autowired)
public class OrderSubscriptionMessageGrpcService extends OrderSubscriptionMessageServiceRpcGrpc.OrderSubscriptionMessageServiceRpcImplBase {
  private final OrderSubscriptionMessageDbService orderSubscriptionMessageDbService;

  @Override
  @DS(DATA_SOURCE_MASTER)
  public void createOrderSubscriptionMessage(final CreateOrderSubscriptionMessageRpcRequest request,
      final StreamObserver<CreateOrderSubscriptionMessageRpcResponse> responseObserver) {
    final var builder = CreateOrderSubscriptionMessageRpcResponse.newBuilder();
    try {
      final var data = orderSubscriptionMessageDbService.createOrderSubscriptionMessage(request);
      builder.setId(data.getId());
    } catch (final Exception e) {
      log.error("createOrderSubscriptionMessage error;", e);
      responseObserver.onError(Status.INTERNAL.withDescription(e.getMessage()).asException());
      return;
    }
    responseObserver.onNext(builder.build());
    responseObserver.onCompleted();
  }

  @Override
  @DS(DATA_SOURCE_MASTER)
  public void getUserOrderSubscriptionMessageHistory(final GetUserOrderSubscriptionMessageHistoryRpcRequest request,
      final StreamObserver<GetUserOrderSubscriptionMessageHistoryRpcResponse> responseObserver) {
    final var builder = GetUserOrderSubscriptionMessageHistoryRpcResponse.newBuilder();
    try {
      final var list = orderSubscriptionMessageDbService.getUserOrderSubscriptionMessageHistory(request);
      builder.addAllResponse(list.stream()
          .map(this::convertToRpcResponse)
          .collect(Collectors.toList()));
    } catch (final Exception e) {
      log.error("getUserOrderSubscriptionMessageHistory error;", e);
      responseObserver.onError(Status.INTERNAL.withDescription(e.getMessage()).asException());
      return;
    }
    responseObserver.onNext(builder.build());
    responseObserver.onCompleted();
  }

  private OrderSubscriptionMessageRpcResponse convertToRpcResponse(OrderSubscriptionMessageDb db) {
    OrderSubscriptionMessageRpcResponse.Builder builder = OrderSubscriptionMessageRpcResponse.newBuilder();
    builder.setId(db.getId());

    if (db.getJumpPath() != null) {
      builder.setJumpPath(db.getJumpPath());
    }

    // 转换LocalDateTime为时间戳(毫秒),与项目中其他时间字段保持一致
    if (db.getOrderPeriodStartDate() != null) {
      long timestamp = db.getOrderPeriodStartDate().atZone(ZoneId.systemDefault())
          .toInstant()
          .toEpochMilli();
      builder.setOrderPeriodStartDate(timestamp);
    }

    if (db.getOrderPeriodEndDate() != null) {
      long timestamp = db.getOrderPeriodEndDate().atZone(ZoneId.systemDefault())
          .toInstant()
          .toEpochMilli();
      builder.setOrderPeriodEndDate(timestamp);
    }

    // 转换LocalDateTime为时间戳(毫秒)
    if (db.getCreatedAt() != null) {
      long timestamp = db.getCreatedAt().atZone(ZoneId.systemDefault())
          .toInstant()
          .toEpochMilli();
      builder.setCreatedAt(timestamp);
    }

    builder.setIsDeleted(db.getIsDeleted() != null && db.getIsDeleted());
    return builder.build();
  }

  @Override
  @DS(DATA_SOURCE_MASTER)
  public void batchDeleteOrderSubscriptionMessagesByIds(final BatchDeleteOrderSubscriptionMessagesByIdsRpcRequest request,
      final StreamObserver<BatchDeleteOrderSubscriptionMessagesByIdsRpcResponse> responseObserver) {
    final var builder = BatchDeleteOrderSubscriptionMessagesByIdsRpcResponse.newBuilder();
    try {
      final int affectedRows = orderSubscriptionMessageDbService.batchDeleteByIds(
          request.getIdsList(),
          request.getEnterpriseId());
      builder.setAffectedRows(affectedRows);
    } catch (final Exception e) {
      log.error("batchDeleteOrderSubscriptionMessagesByIds error;", e);
      responseObserver.onError(Status.INTERNAL.withDescription(e.getMessage()).asException());
      return;
    }
    responseObserver.onNext(builder.build());
    responseObserver.onCompleted();
  }

  @Override
  @DS(DATA_SOURCE_MASTER)
  public void batchPhysicalDeleteOrderSubscriptionMessagesByIds(
      final BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcRequest request,
      final StreamObserver<BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcResponse> responseObserver) {
    final var builder = BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcResponse.newBuilder();
    try {
      final int affectedRows = orderSubscriptionMessageDbService.batchPhysicalDeleteByIds(
          request.getIdsList(),
          request.getEnterpriseId());
      builder.setAffectedRows(affectedRows);
    } catch (final Exception e) {
      log.error("batchPhysicalDeleteOrderSubscriptionMessagesByIds error;", e);
      responseObserver.onError(Status.INTERNAL.withDescription(e.getMessage()).asException());
      return;
    }
    responseObserver.onNext(builder.build());
    responseObserver.onCompleted();
  }
}