diff --git a/pkg/hold/admin/admin.go b/pkg/hold/admin/admin.go index 46e30ec..4d11295 100644 --- a/pkg/hold/admin/admin.go +++ b/pkg/hold/admin/admin.go @@ -423,6 +423,7 @@ func (ui *AdminUI) RegisterRoutes(r chi.Router) { r.Get("/admin/api/stats", ui.handleStatsAPI) r.Get("/admin/api/top-users", ui.handleTopUsersAPI) r.Get("/admin/api/relay/status", ui.handleRelayStatus) + r.Get("/admin/api/crew/member", ui.handleCrewMemberInfo) // Logout r.Post("/admin/auth/logout", ui.handleLogout) diff --git a/pkg/hold/admin/handlers_crew.go b/pkg/hold/admin/handlers_crew.go index 9a73eef..9598637 100644 --- a/pkg/hold/admin/handlers_crew.go +++ b/pkg/hold/admin/handlers_crew.go @@ -9,11 +9,10 @@ import ( "time" "atcr.io/pkg/atproto" - "atcr.io/pkg/hold/pds" "github.com/go-chi/chi/v5" ) -// CrewMemberView represents a crew member for display +// CrewMemberView represents a crew member for display (populated row) type CrewMemberView struct { RKey string DID string @@ -28,6 +27,17 @@ type CrewMemberView struct { AddedAt time.Time } +// CrewSkeletonView is the minimal crew member data for skeleton rendering. +// Contains only data available from the MST walk (no network calls). +type CrewSkeletonView struct { + RKey string + DID string + Role string + Permissions []string + Tier string + AddedAt time.Time +} + // resolveHandle attempts to resolve a DID to a handle // Returns empty string if resolution fails func resolveHandle(ctx context.Context, did string) string { @@ -50,18 +60,14 @@ type TierOption struct { Limit string } -// getCrewViews builds the crew member view list -func (ui *AdminUI) getCrewViews(ctx context.Context) ([]CrewMemberView, error) { - crew, err := ui.pds.ListCrewMembers(ctx) +// handleCrewTab returns the crew tab content (HTMX partial). +// Only does the MST walk — no handle resolution or usage queries. +// Each row lazy-loads its details via handleCrewMemberInfo. +func (ui *AdminUI) handleCrewTab(w http.ResponseWriter, r *http.Request) { + crew, err := ui.pds.ListCrewMembers(r.Context()) if err != nil { - return nil, err - } - - // Single bulk query for all user quotas - allQuotas, err := ui.pds.GetAllUserQuotas(ctx) - if err != nil { - slog.Warn("Failed to get all user quotas for crew views", "error", err) - allQuotas = make(map[string]*pds.QuotaStats) + http.Error(w, "Failed to list crew: "+err.Error(), http.StatusInternalServerError) + return } defaultTier := "default" @@ -69,68 +75,94 @@ func (ui *AdminUI) getCrewViews(ctx context.Context) ([]CrewMemberView, error) { defaultTier = ui.quotaMgr.GetDefaultTier() } - var crewViews []CrewMemberView + var skeletons []CrewSkeletonView for _, member := range crew { tier := member.Record.Tier if tier == "" { tier = defaultTier } - - view := CrewMemberView{ + skeletons = append(skeletons, CrewSkeletonView{ RKey: member.Rkey, DID: member.Record.Member, - Handle: resolveHandle(ctx, member.Record.Member), Role: member.Record.Role, Permissions: member.Record.Permissions, Tier: tier, AddedAt: parseTime(member.Record.AddedAt), - } + }) + } - usage := int64(0) - if q, ok := allQuotas[member.Record.Member]; ok { - usage = q.TotalSize - } + sort.Slice(skeletons, func(i, j int) bool { + return skeletons[i].AddedAt.After(skeletons[j].AddedAt) + }) - if ui.quotaMgr != nil && ui.quotaMgr.IsEnabled() { - if limit := ui.quotaMgr.GetTierLimit(tier); limit != nil { - view.TierLimit = formatHumanBytes(*limit) - if *limit > 0 { - view.UsagePercent = int(float64(usage) / float64(*limit) * 100) - } - } else { - view.TierLimit = "Unlimited" + data := struct { + Crew []CrewSkeletonView + }{ + Crew: skeletons, + } + ui.renderTemplate(w, "partials/tab_crew.html", data) +} + +// handleCrewMemberInfo returns a fully populated crew member row (HTMX partial). +// Called per-row via hx-trigger="load" — resolves handle and fetches usage. +func (ui *AdminUI) handleCrewMemberInfo(w http.ResponseWriter, r *http.Request) { + ctx := r.Context() + rkey := r.URL.Query().Get("rkey") + if rkey == "" { + http.Error(w, "Missing rkey parameter", http.StatusBadRequest) + return + } + + _, member, err := ui.pds.GetCrewMember(ctx, rkey) + if err != nil { + slog.Warn("Failed to get crew member for lazy load", "rkey", rkey, "error", err) + http.Error(w, "Crew member not found", http.StatusNotFound) + return + } + + handle := resolveHandle(ctx, member.Member) + + usage := int64(0) + if q, err := ui.pds.GetQuotaForUser(ctx, member.Member); err == nil { + usage = q.TotalSize + } + + defaultTier := "default" + if ui.quotaMgr != nil && ui.quotaMgr.IsEnabled() { + defaultTier = ui.quotaMgr.GetDefaultTier() + } + + tier := member.Tier + if tier == "" { + tier = defaultTier + } + + view := CrewMemberView{ + RKey: rkey, + DID: member.Member, + Handle: handle, + Role: member.Role, + Permissions: member.Permissions, + Tier: tier, + CurrentUsage: usage, + UsageHuman: formatHumanBytes(usage), + AddedAt: parseTime(member.AddedAt), + } + + if ui.quotaMgr != nil && ui.quotaMgr.IsEnabled() { + if limit := ui.quotaMgr.GetTierLimit(tier); limit != nil { + view.TierLimit = formatHumanBytes(*limit) + if *limit > 0 { + view.UsagePercent = int(float64(usage) / float64(*limit) * 100) } } else { view.TierLimit = "Unlimited" } - - view.CurrentUsage = usage - view.UsageHuman = formatHumanBytes(view.CurrentUsage) - - crewViews = append(crewViews, view) + } else { + view.TierLimit = "Unlimited" } - sort.Slice(crewViews, func(i, j int) bool { - return crewViews[i].CurrentUsage > crewViews[j].CurrentUsage - }) - - return crewViews, nil -} - -// handleCrewTab returns the crew tab content (HTMX partial) -func (ui *AdminUI) handleCrewTab(w http.ResponseWriter, r *http.Request) { - crewViews, err := ui.getCrewViews(r.Context()) - if err != nil { - http.Error(w, "Failed to list crew: "+err.Error(), http.StatusInternalServerError) - return - } - - data := struct { - Crew []CrewMemberView - }{ - Crew: crewViews, - } - ui.renderTemplate(w, "partials/tab_crew.html", data) + ui.renderTemplate(w, "partials/crew_member_row.html", view) } // handleCrewAddForm displays the add crew form diff --git a/pkg/hold/admin/templates/partials/crew_member_row.html b/pkg/hold/admin/templates/partials/crew_member_row.html new file mode 100644 index 0000000..935e43b --- /dev/null +++ b/pkg/hold/admin/templates/partials/crew_member_row.html @@ -0,0 +1,38 @@ +{{define "partials/crew_member_row.html"}} + + +
+ {{if .Handle}}{{.Handle}}
{{end}} + {{.DID}} +
+ + {{.Role}} + + {{range .Permissions}} + {{.}} + {{end}} + + + {{.Tier}} +
{{.TierLimit}} + + +
+ {{.UsageHuman}} + + {{.UsagePercent}}% +
+ + {{formatTime .AddedAt}} + +
+ + {{ icon "pencil" "size-4" }} + + +
+ + +{{end}} diff --git a/pkg/hold/admin/templates/partials/tab_crew.html b/pkg/hold/admin/templates/partials/tab_crew.html index f307833..67f2caf 100644 --- a/pkg/hold/admin/templates/partials/tab_crew.html +++ b/pkg/hold/admin/templates/partials/tab_crew.html @@ -36,11 +36,15 @@ {{range .Crew}} - +
- {{if .Handle}}{{.Handle}}
{{end}} - {{.DID}} + +
+ {{truncate .DID 32}}
{{.Role}} @@ -51,26 +55,12 @@ {{.Tier}} -
{{.TierLimit}} - -
- {{.UsageHuman}} - - {{.UsagePercent}}% -
+ + {{formatTime .AddedAt}} - -
- - {{ icon "pencil" "size-4" }} - - -
- + {{end}} diff --git a/pkg/hold/gc/gc.go b/pkg/hold/gc/gc.go index 4c4e036..af494ba 100644 --- a/pkg/hold/gc/gc.go +++ b/pkg/hold/gc/gc.go @@ -741,33 +741,48 @@ func (gc *GarbageCollector) scanOrphanedBlobDetails(ctx context.Context, referen return orphaned, totalBlobs, nil } +// reconcileBatchSize is the number of layer records per repo commit. +// Batching reduces firehose events from N to N/batchSize. +const reconcileBatchSize = 200 + // reconcileMissingRecords creates layer records for manifest+layer pairs that are missing. +// Records are batched into single commits to avoid flooding relays with firehose events. func (gc *GarbageCollector) reconcileMissingRecords(ctx context.Context, missing []MissingRecordDetail, result *GCResult) { - for i, m := range missing { - record := atproto.NewLayerRecord( - m.Digest, - m.Size, - m.MediaType, - m.UserDID, - m.ManifestURI, - ) - if _, _, err := gc.pds.CreateLayerRecord(ctx, record); err != nil { - gc.logger.Error("Failed to create reconciled layer record", - "digest", m.Digest, - "manifest", m.ManifestURI, + for i := 0; i < len(missing); i += reconcileBatchSize { + end := i + reconcileBatchSize + if end > len(missing) { + end = len(missing) + } + chunk := missing[i:end] + + records := make([]*atproto.LayerRecord, 0, len(chunk)) + for _, m := range chunk { + records = append(records, atproto.NewLayerRecord( + m.Digest, + m.Size, + m.MediaType, + m.UserDID, + m.ManifestURI, + )) + } + + created, err := gc.pds.BatchCreateLayerRecords(ctx, records) + if err != nil { + gc.logger.Error("Failed to create reconciled layer batch", + "batchStart", i, + "batchSize", len(chunk), "error", err) continue } - result.RecordsReconciled++ - if result.RecordsReconciled%100 == 0 { - gc.logger.Info("Reconciliation progress", - "created", result.RecordsReconciled, - "total", len(missing)) - } - // Throttle: ramp delay based on record index to avoid flooding relays - delay := max(10*time.Millisecond, time.Duration(i)*100*time.Microsecond) - time.Sleep(delay) + result.RecordsReconciled += int64(created) + gc.logger.Info("Reconciliation progress", + "created", result.RecordsReconciled, + "total", len(missing), + "batch", len(chunk)) + + // Small delay between batches as a courtesy to relays + time.Sleep(100 * time.Millisecond) } if result.RecordsReconciled > 0 { diff --git a/pkg/hold/pds/layer.go b/pkg/hold/pds/layer.go index 114b3d4..a2f5841 100644 --- a/pkg/hold/pds/layer.go +++ b/pkg/hold/pds/layer.go @@ -3,9 +3,11 @@ package pds import ( "context" "fmt" + "log/slog" "atcr.io/pkg/atproto" "atcr.io/pkg/hold/quota" + indigoatproto "github.com/bluesky-social/indigo/api/atproto" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/repo" ) @@ -41,6 +43,45 @@ func (p *HoldPDS) CreateLayerRecord(ctx context.Context, record *atproto.LayerRe return rkey, recordCID.String(), nil } +// BatchCreateLayerRecords creates multiple layer records in a single repo commit. +// This produces one firehose event instead of one per record. +// Invalid records are skipped with a warning. Returns the number of records created. +func (p *HoldPDS) BatchCreateLayerRecords(ctx context.Context, records []*atproto.LayerRecord) (int, error) { + var writes []*indigoatproto.RepoApplyWrites_Input_Writes_Elem + + for _, record := range records { + if record.Type != atproto.LayerCollection { + slog.Warn("Skipping invalid record type in batch", "type", record.Type) + continue + } + if record.Digest == "" { + slog.Warn("Skipping record with empty digest in batch") + continue + } + if record.Size <= 0 { + slog.Warn("Skipping record with non-positive size in batch", "size", record.Size) + continue + } + + writes = append(writes, &indigoatproto.RepoApplyWrites_Input_Writes_Elem{ + RepoApplyWrites_Create: &indigoatproto.RepoApplyWrites_Create{ + Collection: atproto.LayerCollection, + Value: &lexutil.LexiconTypeDecoder{Val: record}, + }, + }) + } + + if len(writes) == 0 { + return 0, nil + } + + if err := p.repomgr.BatchWrite(ctx, p.uid, writes); err != nil { + return 0, fmt.Errorf("batch write failed: %w", err) + } + + return len(writes), nil +} + // GetLayerRecord retrieves a specific layer record by rkey // Note: This is a simplified implementation. For production, you may need to pass the CID func (p *HoldPDS) GetLayerRecord(ctx context.Context, rkey string) (*atproto.LayerRecord, error) {