package store import ( "context" "fmt" "time" "github.com/jackc/pgx/v5" "github.com/netpulse/netpulse/server/internal/crypto" "github.com/netpulse/netpulse/server/internal/gitstore" ) // GitSyncStat — підсумок переливання історії в Git. type GitSyncStat struct { Total int Written int Unchanged int Skipped int } // SyncGit переливає збережені версії конфігів у Git. // // Потрібна двічі. Перший раз — коли версіювання вмикають на інсталяції, // яка вже місяцями збирає конфіги: без цього репозиторій починається з // наступної зміни, і вся накопичена історія лишається невидимою. // Другий — коли диск із репозиторієм втрачено. Тіла лежать зашифровані // в базі, тож Git тут вторинний і повністю відтворюваний; зворотне // невірно, і саме тому база лишається джерелом істини. // // Версії йдуть у хронологічному порядку — інакше в гілці пристрою // новіший конфіг став би батьком старішого, і історія читалась би // навиворіт. func (s *Store) SyncGit( ctx context.Context, tenantID string, ring *crypto.Keyring, dryRun bool, onProgress func(deviceName, configType string, at time.Time, res gitstore.Result), ) (GitSyncStat, error) { var stat GitSyncStat if s.git == nil { return stat, gitstore.ErrDisabled } type row struct { id string deviceID string deviceName string configType string collected time.Time } var rows []row err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error { r, err := tx.Query(ctx, ` SELECT c.id::text, c.device_id::text, d.name, c.config_type, c.collected_at FROM ncm.configs c JOIN inv.devices d ON d.id = c.device_id WHERE c.tenant_id = $1 ORDER BY c.device_id, c.config_type, c.collected_at `, tenantID) if err != nil { return err } defer r.Close() for r.Next() { var x row if err := r.Scan(&x.id, &x.deviceID, &x.deviceName, &x.configType, &x.collected); err != nil { return err } rows = append(rows, x) } return r.Err() }) if err != nil { return stat, err } stat.Total = len(rows) for _, x := range rows { body, _, err := s.ConfigBody(ctx, tenantID, x.id, ring) if err != nil { // Версія без тіла — не привід зупиняти переливання решти: // одна зіпсована стрічка не має коштувати всієї історії. stat.Skipped++ continue } path := fmt.Sprintf("%s/%s.cfg", sanitizePath(x.deviceName), x.configType) branch := "device/" + x.deviceID if dryRun { stat.Written++ if onProgress != nil { onProgress(x.deviceName, x.configType, x.collected, gitstore.Result{}) } continue } res, err := s.git.Write(gitstore.Commit{ Repo: tenantID + ".git", Branch: branch, Path: path, Body: []byte(body), Message: fmt.Sprintf("%s: %s", x.deviceName, x.configType), Author: "NetPulse", Email: "netpulse@localhost", // Час збору, а не час переливання: інакше вся історія // злипається в одну хвилину й перестає відповідати на // питання «коли це змінилось». When: x.collected, }) if err != nil { return stat, fmt.Errorf("версія %s: %w", x.id, err) } if res.Unchanged { stat.Unchanged++ } else { stat.Written++ } if onProgress != nil { onProgress(x.deviceName, x.configType, x.collected, res) } // Координати в базі мають вести на справжній об'єкт Git, інакше // відновлення виглядатиме зробленим, а посилання лишаться // зламаними. if err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error { _, err := tx.Exec(ctx, ` UPDATE ncm.configs SET commit_sha = $3, blob_sha = $4, branch = $5, path = $6 WHERE tenant_id = $1 AND id = $2 `, tenantID, x.id, res.CommitSHA, res.BlobSHA, branch, path) return err }); err != nil { return stat, err } } return stat, nil } // TenantIDs — усі живі тенанти. Потрібно командам, які працюють над // усією інсталяцією, а не над одним кабінетом. func (s *Store) TenantIDs(ctx context.Context) ([]string, error) { rows, err := s.pool.Query(ctx, ` SELECT id::text FROM core.tenants WHERE deleted_at IS NULL ORDER BY created_at `) if err != nil { return nil, err } defer rows.Close() var out []string for rows.Next() { var id string if err := rows.Scan(&id); err != nil { return nil, err } out = append(out, id) } return out, rows.Err() }