/*
 * Copyright 2017 the original author or authors.
 *
 * Licensed under the Apache License, Version 2.0 (the "License");
 * you may not use this file except in compliance with the License.
 * You may obtain a copy of the License at
 *
 *      http://www.apache.org/licenses/LICENSE-2.0
 *
 * Unless required by applicable law or agreed to in writing, software
 * distributed under the License is distributed on an "AS IS" BASIS,
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
 * See the License for the specific language governing permissions and
 * limitations under the License.
 */
package<%=packageName%>;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.Input;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.cloud.stream.messaging.Source;
import org.springframework.messaging.SubscribableChannel;

@EnableBinding({ Source.class, KafkaLoggingSource.Sinks.class })
public class KafkaLoggingSource {

    private static Logger logger = LoggerFactory.getLogger(KafkaLoggingSource.class);

    @StreamListener("input")
    public void log(String message) {
        logger.info("Finished processing {}", message);
    }

    @StreamListener("loggingErr")
    public void errorPersons(String error) {
        logger.warn("ERRORS:{}", error);
    }

    public interface Sinks {

        @Input
        SubscribableChannel loggingErr();
    }
}
