package storage import ( "database/sql" "os" "sync" _ "github.com/mattn/go-sqlite3" ) // Repository defines the interface for telegram persistence. type Repository interface { Insert(telegram string) error LoadUnprocessed() ([]Telegram, error) MarkProcessed(id int64) error Close() error } // Telegram represents a stored telegram record. type Telegram struct { ID int64 Text string } // Store implements Repository using SQLite. type Store struct { mu sync.Mutex dbFile string db *sql.DB initTable bool } const ( driverName = "sqlite3" tableDDL = ` create table IF NOT EXISTS telegram ( [tele_id] INTEGER PRIMARY KEY AUTOINCREMENT, [tele_recv_time] TIMESTAMP NOT NULL DEFAULT (datetime('now', 'localtime')), [tele_processed] int(1) NOT NULL DEFAULT 0, [tele_text] TEXT NOT NULL ) ` insertSQL = "insert into telegram (tele_text) values (?)" countSQL = "select count(*) from telegram where tele_processed=0" loadSQL = "select tele_id, tele_text from telegram where tele_processed=0 Limit 100" updateSQL = "update telegram set tele_processed = 1 where tele_id=?" ) // New opens or creates a SQLite store. func New(dbFile string, init bool) (*Store, error) { if init { os.Remove(dbFile) } db, err := sql.Open(driverName, dbFile) if err != nil { return nil, err } if _, err := db.Exec(tableDDL); err != nil { db.Close() return nil, err } return &Store{ dbFile: dbFile, db: db, initTable: init, }, nil } // Insert saves a telegram to the database. func (s *Store) Insert(telegram string) error { s.mu.Lock() defer s.mu.Unlock() stmt, err := s.db.Prepare(insertSQL) if err != nil { return err } defer stmt.Close() _, err = stmt.Exec(telegram) return err } // LoadUnprocessed returns up to 100 unprocessed telegrams. func (s *Store) LoadUnprocessed() ([]Telegram, error) { rows, err := s.db.Query(loadSQL) if err != nil { return nil, err } defer rows.Close() var telegrams []Telegram for rows.Next() { var t Telegram if err := rows.Scan(&t.ID, &t.Text); err != nil { return nil, err } telegrams = append(telegrams, t) } return telegrams, nil } // MarkProcessed marks a telegram as processed. func (s *Store) MarkProcessed(id int64) error { s.mu.Lock() defer s.mu.Unlock() stmt, err := s.db.Prepare(updateSQL) if err != nil { return err } defer stmt.Close() _, err = stmt.Exec(id) return err } // CountUnprocessed returns the number of unprocessed telegrams. func (s *Store) CountUnprocessed() (int64, error) { var count int64 err := s.db.QueryRow(countSQL).Scan(&count) return count, err } // Close closes the database connection. func (s *Store) Close() error { return s.db.Close() }