Start HTTP before the retention purge so sign-in is not blocked
The first-run purge of network-wide follows can take a long time. Serving the API first lets OAuth persist a session; jetstream and fetch wait until the purge finishes so they do not refill the volume mid-pass.
e861f3e parent: 0e0f3d1 modified
main.go +40 -40 | @@ -119,52 +119,52 @@ func main() { | ||
| 119 | 119 | // unbounded. |
| 120 | 120 | jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger, dbs.CursorStore(), dbs.Users.UserDIDList) |
| 121 | 121 | |
| 122 | - // Purge and reclaim disk space before the streaming jobs start writing | |
| 123 | - // again: once the volume is full every SQLite write fails, including | |
| 124 | - // OAuth sign-in. | |
| 125 | - retentionCtx, cancelRetention := context.WithTimeout(context.Background(), 30*time.Minute) | |
| 126 | - start := time.Now() | |
| 127 | - unknownStats, err := dbs.PurgeUnknownUserRows(retentionCtx) | |
| 128 | - if err != nil { | |
| 129 | - logger.Error("initial purge of unknown-user rows failed", "error", err) | |
| 130 | - } else if unknownStats.Total() > 0 { | |
| 131 | - logger.Info("purged rows of unknown users", "stats", unknownStats) | |
| 132 | - } | |
| 133 | - expired, err := dbs.PurgeExpiredArticles(retentionCtx, *articleRetentionDays) | |
| 134 | - if err != nil { | |
| 135 | - logger.Error("initial article retention purge failed", "error", err) | |
| 136 | - } else if expired > 0 { | |
| 137 | - logger.Info("purged expired articles", "count", expired) | |
| 138 | - } | |
| 139 | - if unknownStats.Total()+expired > 0 { | |
| 140 | - if err := dbs.ReclaimSpace(retentionCtx); err != nil { | |
| 141 | - logger.Error("reclaiming database space incomplete", "error", err) | |
| 142 | - } else { | |
| 143 | - logger.Info("reclaimed database space") | |
| 144 | - } | |
| 145 | - } | |
| 146 | - cancelRetention() | |
| 147 | - logger.Info("initial retention complete", "elapsed", time.Since(start).Round(time.Second)) | |
| 148 | 122 | ctx, cancel := context.WithCancel(context.Background()) |
| 149 | 123 | defer cancel() |
| 150 | 124 | |
| 125 | + // Serve HTTP first so sign-in works while the one-time purge is still | |
| 126 | + // deleting network-wide rows. Jobs that write wait until that pass finishes. | |
| 151 | 127 | go func() { |
| 152 | - if err := scheduler.Run(ctx); err != nil && ctx.Err() == nil { | |
| 153 | - logger.Error("scheduler error", "error", err) | |
| 128 | + retentionCtx, cancelRetention := context.WithTimeout(ctx, 30*time.Minute) | |
| 129 | + defer cancelRetention() | |
| 130 | + start := time.Now() | |
| 131 | + unknownStats, err := dbs.PurgeUnknownUserRows(retentionCtx) | |
| 132 | + if err != nil { | |
| 133 | + logger.Error("initial purge of unknown-user rows failed", "error", err) | |
| 134 | + } else if unknownStats.Total() > 0 { | |
| 135 | + logger.Info("purged rows of unknown users", "stats", unknownStats) | |
| 154 | 136 | } |
| 155 | - }() | |
| 156 | - go func() { | |
| 157 | - if err := cron.Run(ctx); err != nil && ctx.Err() == nil { | |
| 158 | - logger.Error("cron error", "error", err) | |
| 137 | + expired, err := dbs.PurgeExpiredArticles(retentionCtx, *articleRetentionDays) | |
| 138 | + if err != nil { | |
| 139 | + logger.Error("initial article retention purge failed", "error", err) | |
| 140 | + } else if expired > 0 { | |
| 141 | + logger.Info("purged expired articles", "count", expired) | |
| 159 | 142 | } |
| 160 | - }() | |
| 161 | - go func() { | |
| 162 | - srv.PeriodicSync(ctx, *syncInterval) | |
| 163 | - }() | |
| 164 | - go func() { | |
| 165 | - srv.BackfillFromCollectionDir(ctx, *collectionDirURL, *backfillConcurrency) | |
| 166 | - }() | |
| 167 | - go func() { | |
| 143 | + if unknownStats.Total()+expired > 0 { | |
| 144 | + if err := dbs.ReclaimSpace(retentionCtx); err != nil { | |
| 145 | + logger.Error("reclaiming database space incomplete", "error", err) | |
| 146 | + } else { | |
| 147 | + logger.Info("reclaimed database space") | |
| 148 | + } | |
| 149 | + } | |
| 150 | + logger.Info("initial retention complete", "elapsed", time.Since(start).Round(time.Second)) | |
| 151 | + | |
| 152 | + go func() { | |
| 153 | + if err := scheduler.Run(ctx); err != nil && ctx.Err() == nil { | |
| 154 | + logger.Error("scheduler error", "error", err) | |
| 155 | + } | |
| 156 | + }() | |
| 157 | + go func() { | |
| 158 | + if err := cron.Run(ctx); err != nil && ctx.Err() == nil { | |
| 159 | + logger.Error("cron error", "error", err) | |
| 160 | + } | |
| 161 | + }() | |
| 162 | + go func() { | |
| 163 | + srv.PeriodicSync(ctx, *syncInterval) | |
| 164 | + }() | |
| 165 | + go func() { | |
| 166 | + srv.BackfillFromCollectionDir(ctx, *collectionDirURL, *backfillConcurrency) | |
| 167 | + }() | |
| 168 | 168 | if err := jetstream.Start(ctx); err != nil && ctx.Err() == nil { |
| 169 | 169 | logger.Error("jetstream error", "error", err) |
| 170 | 170 | } |
| @@ -119,52 +119,52 @@ func main() { | |||
| 119 | // unbounded. | 119 | // unbounded. |
| 120 | jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger, dbs.CursorStore(), dbs.Users.UserDIDList) | 120 | jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger, dbs.CursorStore(), dbs.Users.UserDIDList) |
| 121 | 121 | ||
| 122 | - // Purge and reclaim disk space before the streaming jobs start writing | ||
| 123 | - // again: once the volume is full every SQLite write fails, including | ||
| 124 | - // OAuth sign-in. | ||
| 125 | - retentionCtx, cancelRetention := context.WithTimeout(context.Background(), 30*time.Minute) | ||
| 126 | - start := time.Now() | ||
| 127 | - unknownStats, err := dbs.PurgeUnknownUserRows(retentionCtx) | ||
| 128 | - if err != nil { | ||
| 129 | - logger.Error("initial purge of unknown-user rows failed", "error", err) | ||
| 130 | - } else if unknownStats.Total() > 0 { | ||
| 131 | - logger.Info("purged rows of unknown users", "stats", unknownStats) | ||
| 132 | - } | ||
| 133 | - expired, err := dbs.PurgeExpiredArticles(retentionCtx, *articleRetentionDays) | ||
| 134 | - if err != nil { | ||
| 135 | - logger.Error("initial article retention purge failed", "error", err) | ||
| 136 | - } else if expired > 0 { | ||
| 137 | - logger.Info("purged expired articles", "count", expired) | ||
| 138 | - } | ||
| 139 | - if unknownStats.Total()+expired > 0 { | ||
| 140 | - if err := dbs.ReclaimSpace(retentionCtx); err != nil { | ||
| 141 | - logger.Error("reclaiming database space incomplete", "error", err) | ||
| 142 | - } else { | ||
| 143 | - logger.Info("reclaimed database space") | ||
| 144 | - } | ||
| 145 | - } | ||
| 146 | - cancelRetention() | ||
| 147 | - logger.Info("initial retention complete", "elapsed", time.Since(start).Round(time.Second)) | ||
| 148 | ctx, cancel := context.WithCancel(context.Background()) | 122 | ctx, cancel := context.WithCancel(context.Background()) |
| 149 | defer cancel() | 123 | defer cancel() |
| 150 | 124 | ||
| 125 | + // Serve HTTP first so sign-in works while the one-time purge is still | ||
| 126 | + // deleting network-wide rows. Jobs that write wait until that pass finishes. | ||
| 151 | go func() { | 127 | go func() { |
| 152 | - if err := scheduler.Run(ctx); err != nil && ctx.Err() == nil { | 128 | + retentionCtx, cancelRetention := context.WithTimeout(ctx, 30*time.Minute) |
| 153 | - logger.Error("scheduler error", "error", err) | 129 | + defer cancelRetention() |
| 130 | + start := time.Now() | ||
| 131 | + unknownStats, err := dbs.PurgeUnknownUserRows(retentionCtx) | ||
| 132 | + if err != nil { | ||
| 133 | + logger.Error("initial purge of unknown-user rows failed", "error", err) | ||
| 134 | + } else if unknownStats.Total() > 0 { | ||
| 135 | + logger.Info("purged rows of unknown users", "stats", unknownStats) | ||
| 154 | } | 136 | } |
| 155 | - }() | 137 | + expired, err := dbs.PurgeExpiredArticles(retentionCtx, *articleRetentionDays) |
| 156 | - go func() { | 138 | + if err != nil { |
| 157 | - if err := cron.Run(ctx); err != nil && ctx.Err() == nil { | 139 | + logger.Error("initial article retention purge failed", "error", err) |
| 158 | - logger.Error("cron error", "error", err) | 140 | + } else if expired > 0 { |
| 141 | + logger.Info("purged expired articles", "count", expired) | ||
| 159 | } | 142 | } |
| 160 | - }() | 143 | + if unknownStats.Total()+expired > 0 { |
| 161 | - go func() { | 144 | + if err := dbs.ReclaimSpace(retentionCtx); err != nil { |
| 162 | - srv.PeriodicSync(ctx, *syncInterval) | 145 | + logger.Error("reclaiming database space incomplete", "error", err) |
| 163 | - }() | 146 | + } else { |
| 164 | - go func() { | 147 | + logger.Info("reclaimed database space") |
| 165 | - srv.BackfillFromCollectionDir(ctx, *collectionDirURL, *backfillConcurrency) | 148 | + } |
| 166 | - }() | 149 | + } |
| 167 | - go func() { | 150 | + logger.Info("initial retention complete", "elapsed", time.Since(start).Round(time.Second)) |
| 151 | + | ||
| 152 | + go func() { | ||
| 153 | + if err := scheduler.Run(ctx); err != nil && ctx.Err() == nil { | ||
| 154 | + logger.Error("scheduler error", "error", err) | ||
| 155 | + } | ||
| 156 | + }() | ||
| 157 | + go func() { | ||
| 158 | + if err := cron.Run(ctx); err != nil && ctx.Err() == nil { | ||
| 159 | + logger.Error("cron error", "error", err) | ||
| 160 | + } | ||
| 161 | + }() | ||
| 162 | + go func() { | ||
| 163 | + srv.PeriodicSync(ctx, *syncInterval) | ||
| 164 | + }() | ||
| 165 | + go func() { | ||
| 166 | + srv.BackfillFromCollectionDir(ctx, *collectionDirURL, *backfillConcurrency) | ||
| 167 | + }() | ||
| 168 | if err := jetstream.Start(ctx); err != nil && ctx.Err() == nil { | 168 | if err := jetstream.Start(ctx); err != nil && ctx.Err() == nil { |
| 169 | logger.Error("jetstream error", "error", err) | 169 | logger.Error("jetstream error", "error", err) |
| 170 | } | 170 | } |