/*
Nekünk át kell adni adatokat egy forrásból egy fogyasztónak.
A forrás kis csomagokban (~tíz rekord) adja az adatokat, míg a fogyasztó hatékonyabban működik nagyobb csomagokkal (~ezer rekord).
Egy valódi példa a Kafka típusú sorokból származó adatok átadása a Clickhouse adatbázisba.
Forrás:
- Szinte végtelen.
- A forrás soha nem ad vissza több mint MaxItems rekordot egy Next hívásban.
- Egy "munkamenet" (egy Pipe függvényhívás) során a forrás minden Next hívásnál új adatokat ad vissza.
- Azonban, újraindítás után a forrás a korábbi "megerősített" pozícióból indul, amit a cookie határoz meg.
Ezért, minden cookie érték, amit a Next visszaad, az adatok tárolása után,
a Commit hívással kell megerősíteni, ugyanabban a sorrendben, ahogyan a Next visszaadta.
Fogadó:
- Egyszerre nem dolgozhat több mint MaxItems-t.
A feladat: implementálni a func Pipe(p Producer, c Consumer) error függvényt,
ami a forrásból olvas, az adatokat MaxItems méretű bufferbe gyűjti, és a fogadóba menti,
azután megerősíti a forrást a haladásról.
*/
const MaxItems = 9999
type Producer interface {
// Next:
// - egy csomag elemet ad vissza
// - cookie, amit a feldolgozás befejezése után kell megerősíteni
// - hiba
Next() (items []any, cookie int, err error)
// Commit a csomag feldolgozottságának jelölésére szolgál
Commit(cookie int) error
}
type Consumer interface {
Process(items []any) error
}
func Pipe(p Producer, c Consumer) error {
// TODO
}