/*
Трябва да прехвърлим данни от източник към потребител.
Източникът предоставя данни на малки пакети (~десет записа), докато потребителят работи по-ефективно с големи пакети (~хиляда записа).
Реален пример е прехвърлянето на данни от опашки тип Kafka към база данни Clickhouse.
Източник:
- Практически безкраен.
- Източникът никога не връща повече от MaxItems записи в един обаждане на Next.
- В рамките на една "сесия" (едно обаждане на функцията Pipe), източникът връща нови данни при всяко Next.
- Въпреки това, след рестарт, източникът започва от предишната "потвърдена" позиция, определена от cookie.
Затова, *всяка* стойност на cookie, която Next връща, след като данните са запазени в приемника,
трябва да бъде потвърдена с извикване на Commit, в същия ред, в който са били върнати от Next.
Приемник:
- Не може да обработи повече от MaxItems наведнъж.
Трябва да се реализира функцията func Pipe(p Producer, c Consumer) error,
която чете данни от източника, ги групира в буфер с размер не по-голям от MaxItems и ги запазва в приемника,
след което потвърждава напредъка в източника.
*/
const MaxItems = 9999
type Producer interface {
// Next връща:
// - пакет от елементи за обработка
// - cookie за потвърждение, когато обработката е завършена
// - грешка
Next() (items []any, cookie int, err error)
// Commit се използва за маркиране на пакет данни като обработен
Commit(cookie int) error
}
type Consumer interface {
Process(items []any) error
}
func Pipe(p Producer, c Consumer) error {
// TODO
}