package com.orchard;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.orchard.models.GlobalParticipantStreams;
import com.orchard.models.OrchardModel;
import io.confluent.ksql.function.udf.UdfParameter;
import io.confluent.ksql.function.udf.UdfSchemaProvider;
import io.confluent.ksql.function.udtf.Udtf;
import io.confluent.ksql.function.udtf.UdtfDescription;
import io.confluent.ksql.schema.ksql.SchemaConverters;
import io.confluent.ksql.schema.ksql.types.SqlType;
import java.util.Arrays;
import java.util.List;
import java.util.stream.Collectors;
import org.apache.kafka.connect.data.Struct;


@UdtfDescription(name = "split_global_participants", description = "Split global participants data by new lines.")
public class SplitGlobalParticipantsUdtf {
    private final String DELIMITER = "\\R";
    private final ObjectMapper mapper = new ObjectMapper();
    private final SchemaConverters.ConnectToSqlTypeConverter converter = SchemaConverters.connectToSqlConverter();

    @Udtf(description = "Splits a multi-line string into a list of structs.", schemaProvider = "provideSchema")
    public List<Struct> split_global_participants(
            @UdfParameter(value = "input", description = "the given string to split") final String input) {
        return Arrays.stream(input.split(DELIMITER)).map(
                item -> {
                    try {
                        OrchardModel gp = mapper.readValue(item, GlobalParticipantStreams.class);
                        return gp.ToKafkaStruct();
                    } catch (JsonProcessingException e) {
                        System.out.println(e.getMessage());
                    };
                    return null;
                }
        ).collect(Collectors.toList());
    }

    @UdfSchemaProvider
    public SqlType provideSchema(final List<SqlType> params) {
        return converter.toSqlType(GlobalParticipantStreams.KafkaSchema);
    }
}
