package connection import ( "context" "log/slog" "github.com/bluesky-social/indigo/atproto/auth/oauth" ) // DrainResult summarizes what happened during a queue drain. type DrainResult struct { // Written is the count of pending items successfully flushed to the // authenticated user's PDS. Written int // Skipped is the count of pending items that hit a write error and were // left in the queue (so a future drain can retry). Skipped int } // DrainHook is called after each successfully drained pending connection. // It gives the caller a chance to perform side-effects like auto-check-in // without coupling the connection package to the checkin/badge packages. type DrainHook func(ctx context.Context, sess *oauth.ClientSession, item PendingItem) // Drain flushes every pending_connection row for sess.AccountDID by calling // Put for each, then deleting the row on success. Errors on individual rows // are logged (if logger is non-nil) and the row is left in the queue for a // future retry. // // pdsHost is the PDS URL of the draining user, used to dedup against existing // connection records on their PDS. // // If onDrain is non-nil it is called after each successfully written + // deleted pending item, e.g. to auto-check-in the user to the event. // // Returns a DrainResult summarizing the operation; never returns an error // from the per-row writes — the goal is best-effort flush on login, not // blocking the user's redirect to /profile. // // A wrapper-level error (e.g. queue lookup failure) is returned as-is. func Drain(ctx context.Context, q *Queue, sess *oauth.ClientSession, pdsHost string, logger *slog.Logger, onDrain DrainHook) (DrainResult, error) { res := DrainResult{} if sess == nil { return res, errNoSession } target := sess.Data.AccountDID items, err := q.List(ctx, target, 0) if err != nil { return res, err } for _, item := range items { // Dedup: skip if the user already has this connection on their PDS. if pdsHost != "" && HasConnection(ctx, pdsHost, target, item.InitiatorDID, item.EventURI) { if logger != nil { logger.Debug("connection drain: duplicate skipped", "target", target.String(), "initiator", item.InitiatorDID.String(), "event", item.EventURI, ) } // Write bookmark for this known connection if q.db != nil { viewerDID := target.String() _, _ = q.db.ExecContext(ctx, ` INSERT INTO connection_notes (viewer_did, target_did, notes, follow_up, updated_at) VALUES (?, ?, '', 0, CURRENT_TIMESTAMP) ON CONFLICT(viewer_did, target_did) DO NOTHING `, viewerDID, item.InitiatorDID.String()) } // Delete the queue row — no need to retry. _ = q.Delete(ctx, item.ID) res.Skipped++ // Still fire the hook so auto-checkin happens even for skipped connections. if onDrain != nil { onDrain(ctx, sess, item) } continue } _, _, err := Put(ctx, sess, q.db, Record{With: item.InitiatorDID, EventURI: item.EventURI}) if err != nil { if logger != nil { logger.Warn("connection drain: write failed", "target", target.String(), "initiator", item.InitiatorDID.String(), "err", err, ) } res.Skipped++ continue } if err := q.Delete(ctx, item.ID); err != nil { // Wrote the record successfully but couldn't delete the queue // row — next drain will retry and create a duplicate record. Log // loudly so we can spot it. if logger != nil { logger.Error("connection drain: row delete failed; will duplicate on retry", "id", item.ID, "err", err, ) } res.Skipped++ continue } res.Written++ if onDrain != nil { onDrain(ctx, sess, item) } } return res, nil } // errNoSession is returned by Drain when invoked without a session. Kept // package-private — callers should ensure they have a session before calling. var errNoSession = errSentinel("connection drain: nil session") type errSentinel string func (e errSentinel) Error() string { return string(e) }