package com.orchard;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.orchard.models.GlobalSoundRecordingsStreams;
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_sound_recordings", description = "Split global soundrecordings data by new lines.")
public class SplitGlobalSoundRecordingsUdtf {
    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_sound_recordings(
            @UdfParameter(value = "input", description = "the given string to split") final String input) {
        return Arrays.stream(input.split(DELIMITER)).map(
                item -> {
                    GlobalSoundRecordingsStreams gs = null;
                    try {
                        gs = mapper.readValue(item, GlobalSoundRecordingsStreams.class);
                        return gs.ToKafkaStruct();
                    } catch (JsonProcessingException e) {
                        System.out.println(e.getMessage());
                    };
                    return null;
                }
        ).filter(item -> item != null).collect(Collectors.toList());
    }

    @UdfSchemaProvider
    public SqlType provideSchema(final List<SqlType> params) {
        SchemaConverters.ConnectToSqlTypeConverter converter = SchemaConverters.connectToSqlConverter();
        return converter.toSqlType(GlobalSoundRecordingsStreams.KafkaSchema);
    }
}
