diff --git a/internal/repository/postgres.go b/internal/repository/postgres.go index 3c1e31e..df7523f 100644 --- a/internal/repository/postgres.go +++ b/internal/repository/postgres.go @@ -6,7 +6,8 @@ import ( "encoding/json" "errors" "fmt" - "sort" + "os" + "runtime" "strconv" "strings" "time" @@ -18,6 +19,36 @@ import ( "github.com/jackc/pgx/v5/pgxpool" ) +// #region agent log +func agentDebugNDJSON3214(hypothesisID, location, message string, data map[string]any) { + if os.Getenv("EVOBGP_DEBUG_LOG") != "1" { + return + } + f, err := os.OpenFile("debug-3214dc.log", os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) + if err != nil { + return + } + defer f.Close() + var ms runtime.MemStats + runtime.ReadMemStats(&ms) + payload := map[string]any{ + "sessionId": "3214dc", + "hypothesisId": hypothesisID, + "location": location, + "message": message, + "data": data, + "timestamp": time.Now().UnixMilli(), + "allocBytes": ms.Alloc, + } + b, err := json.Marshal(payload) + if err != nil { + return + } + _, _ = f.Write(append(b, '\n')) +} + +// #endregion + // Postgres implements store.Backend using pgxpool. type Postgres struct { pool *pgxpool.Pool @@ -42,22 +73,15 @@ func (p *Postgres) DemoIDs() (tenant, moduleCDN, moduleIP, revision, speaker str func (p *Postgres) MaterializedPrefixStats() (max int, sum int) { ctx := context.Background() - rows, err := p.pool.Query(ctx, ` - SELECT COALESCE((meta_json->>'materialized_prefix_count')::int, 0) AS n - FROM config_revision`) + // Агрегация в БД — не тащим все строки config_revision в память. + err := p.pool.QueryRow(ctx, ` + SELECT + COALESCE(MAX((meta_json->>'materialized_prefix_count')::int), 0), + COALESCE(SUM((meta_json->>'materialized_prefix_count')::int), 0) + FROM config_revision`).Scan(&max, &sum) if err != nil { return 0, 0 } - defer rows.Close() - for rows.Next() { - var n int - if rows.Scan(&n) == nil { - sum += n - if n > max { - max = n - } - } - } return max, sum } @@ -554,14 +578,29 @@ func (p *Postgres) ListRevisions(tenantID, moduleID string, cursor string, limit if limit <= 0 { limit = 50 } + off := 0 + if cursor != "" { + if n, err := strconv.Atoi(cursor); err == nil && n >= 0 { + off = n + } + } ctx := context.Background() - q := `SELECT id::text, module_id::text, content_hash, parent_revision_id::text, meta_json, created_at FROM config_revision WHERE tenant_id=$1` + // List endpoint only needs materialized_prefix_count from meta_json — not full preview_fragments blobs. + // LIMIT/OFFSET in SQL avoids loading every revision for the tenant into memory (was O(N) per request). + q := `SELECT id::text, module_id::text, content_hash, parent_revision_id::text, + COALESCE((meta_json->>'materialized_prefix_count')::int, 0), created_at + FROM config_revision WHERE tenant_id=$1` args := []any{tenantID} + n := 2 if moduleID != "" { - q += ` AND module_id = $2` + q += fmt.Sprintf(` AND module_id=$%d`, n) args = append(args, moduleID) + n++ } q += ` ORDER BY created_at DESC` + // Fetch limit+1 rows to compute has_more without COUNT(*). + q += fmt.Sprintf(` LIMIT $%d OFFSET $%d`, n, n+1) + args = append(args, limit+1, off) rows, err := p.pool.Query(ctx, q, args...) if err != nil { return nil, "", false @@ -572,58 +611,53 @@ func (p *Postgres) ListRevisions(tenantID, moduleID string, cursor string, limit var r store.Revision r.TenantID = tenantID var mod, parent *string - var meta []byte - if err := rows.Scan(&r.ID, &mod, &r.ContentHash, &parent, &meta, &r.CreatedAt); err != nil { + var mpc int + if err := rows.Scan(&r.ID, &mod, &r.ContentHash, &parent, &mpc, &r.CreatedAt); err != nil { continue } if mod != nil { r.ModuleID = *mod } r.ParentRevisionID = strOrNil(parent) - var mj struct { - PreviewFragments map[string]string `json:"preview_fragments"` - MaterializedPrefixCount int `json:"materialized_prefix_count"` - } - _ = json.Unmarshal(meta, &mj) - if mj.PreviewFragments == nil { - mj.PreviewFragments = map[string]string{} - } - r.PreviewFragments = mj.PreviewFragments - r.MaterializedPrefixCount = mj.MaterializedPrefixCount + r.MaterializedPrefixCount = mpc + r.PreviewFragments = map[string]string{} all = append(all, &r) } + agentDebugNDJSON3214("A", "repository/postgres.go:ListRevisions", "list_revisions_fetched", map[string]any{ + "rows": len(all), "limit": limit, "offset": off, + }) + hasMore := len(all) > limit + if hasMore { + all = all[:limit] + } + next := "" + if hasMore { + next = fmt.Sprintf("%d", off+limit) + } + if len(all) == 0 { + return nil, "", false + } + return all, next, hasMore +} + +func (p *Postgres) ListRevisionPrefixes(tenantID, revisionID string, cursor string, limit int) ([]store.PrefixRow, string, bool) { + if limit <= 0 { + limit = 50 + } off := 0 if cursor != "" { if n, err := strconv.Atoi(cursor); err == nil && n >= 0 { off = n } } - end := off + limit - next := "" - hasMore := false - if end > len(all) { - end = len(all) - } else { - hasMore = true - next = fmt.Sprintf("%d", end) - } - if off >= len(all) { - return nil, "", false - } - return all[off:end], next, hasMore -} - -func (p *Postgres) ListRevisionPrefixes(tenantID, revisionID string, cursor string, limit int) ([]store.PrefixRow, string, bool) { - if limit <= 0 { - limit = 50 - } ctx := context.Background() if _, err := p.GetRevision(tenantID, revisionID); err != nil { return nil, "", false } rows, err := p.pool.Query(ctx, ` SELECT prefix::text, community_id::text, source FROM revision_materialized_prefix - WHERE revision_id=$1 ORDER BY id`, revisionID) + WHERE revision_id=$1 ORDER BY id + LIMIT $2 OFFSET $3`, revisionID, limit+1, off) if err != nil { return nil, "", false } @@ -638,25 +672,21 @@ func (p *Postgres) ListRevisionPrefixes(tenantID, revisionID string, cursor stri pr.CommunityID = comm all = append(all, pr) } - off := 0 - if cursor != "" { - if n, err := strconv.Atoi(cursor); err == nil && n >= 0 { - off = n - } + agentDebugNDJSON3214("B", "repository/postgres.go:ListRevisionPrefixes", "list_prefixes_fetched", map[string]any{ + "rows": len(all), "limit": limit, "offset": off, + }) + more := len(all) > limit + if more { + all = all[:limit] } - end := off + limit next := "" - more := false - if end > len(all) { - end = len(all) - } else { - more = true - next = fmt.Sprintf("%d", end) + if more { + next = fmt.Sprintf("%d", off+limit) } - if off >= len(all) { + if len(all) == 0 { return nil, "", false } - return all[off:end], next, more + return all, next, more } func (p *Postgres) CreateRollbackRevision(tenantID, sourceRevisionID string) (string, error) { @@ -697,32 +727,54 @@ func (p *Postgres) RevisionDiff(tenantID, aID, bID string) (map[string]any, erro if _, err := p.GetRevision(tenantID, bID); err != nil { return nil, err } - pa, _, _ := p.ListRevisionPrefixes(tenantID, aID, "", 100000) - pb, _, _ := p.ListRevisionPrefixes(tenantID, bID, "", 100000) - setA := make(map[string]struct{}) - setB := make(map[string]struct{}) - for _, x := range pa { - setA[x.Prefix] = struct{}{} + ctx := context.Background() + var unchanged int + err := p.pool.QueryRow(ctx, ` + SELECT COUNT(*)::int FROM ( + SELECT prefix FROM revision_materialized_prefix WHERE revision_id=$1::uuid + INTERSECT + SELECT prefix FROM revision_materialized_prefix WHERE revision_id=$2::uuid + ) t`, aID, bID).Scan(&unchanged) + if err != nil { + return nil, err } - for _, x := range pb { - setB[x.Prefix] = struct{}{} + // added: в B, нет в A; removed: в A, нет в B — без загрузки полных снапшотов в память. + rowsAdded, err := p.pool.Query(ctx, ` + SELECT prefix::text FROM ( + SELECT prefix FROM revision_materialized_prefix WHERE revision_id=$1::uuid + EXCEPT + SELECT prefix FROM revision_materialized_prefix WHERE revision_id=$2::uuid + ) s ORDER BY 1`, bID, aID) + if err != nil { + return nil, err } - var added, removed []string - unchanged := 0 - for pfx := range setB { - if _, ok := setA[pfx]; !ok { - added = append(added, pfx) - } else { - unchanged++ + defer rowsAdded.Close() + var added []string + for rowsAdded.Next() { + var s string + if err := rowsAdded.Scan(&s); err != nil { + continue } + added = append(added, s) } - for pfx := range setA { - if _, ok := setB[pfx]; !ok { - removed = append(removed, pfx) + rowsRem, err := p.pool.Query(ctx, ` + SELECT prefix::text FROM ( + SELECT prefix FROM revision_materialized_prefix WHERE revision_id=$1::uuid + EXCEPT + SELECT prefix FROM revision_materialized_prefix WHERE revision_id=$2::uuid + ) s ORDER BY 1`, aID, bID) + if err != nil { + return nil, err + } + defer rowsRem.Close() + var removed []string + for rowsRem.Next() { + var s string + if err := rowsRem.Scan(&s); err != nil { + continue } + removed = append(removed, s) } - sort.Strings(added) - sort.Strings(removed) return map[string]any{ "revision_a": aID, "revision_b": bID,