From dfbf789afb039850919cfc7f22eb7f4a9c02e400 Mon Sep 17 00:00:00 2001 From: Aly Raffauf Date: Fri, 3 Jul 2026 16:10:37 -0400 Subject: [PATCH] add getManyToManyCounts and getManyToMany --- README.md | 38 ++++++++- internal/api/many_to_many.go | 79 +++++++++++++++++ internal/api/many_to_many_counts.go | 73 ++++++++++++++++ internal/api/server.go | 2 + internal/api/source.go | 17 +++- internal/store/query.go | 127 ++++++++++++++++++++++++++++ 6 files changed, 330 insertions(+), 6 deletions(-) create mode 100644 internal/api/many_to_many.go create mode 100644 internal/api/many_to_many_counts.go diff --git a/README.md b/README.md index 06bc5c9..7ee9bd7 100644 --- a/README.md +++ b/README.md @@ -43,7 +43,7 @@ go run ./cmd/asterism/ ## API -Asterism implements three endpoints from the [microcosm links XRPC namespace](https://constellation.microcosm.blue/): +Asterism implements all five current endpoints from the [microcosm links XRPC namespace](https://constellation.microcosm.blue/) (the older `/links/*` REST endpoints are deprecated upstream in favor of these and aren't implemented here): ### `GET /xrpc/blue.microcosm.links.getBacklinksCount` @@ -87,11 +87,45 @@ Response: `{"total": 42, "records": [{"did": "...", "collection": "...", "rkey": Records identify the linking record by DID, collection, and rkey. Clients must hydrate display data separately (via AppView, PDS, etc.). +### `GET /xrpc/blue.microcosm.links.getManyToMany` + +Join records linking to a subject with a second field path on those same records — a one-hop join in a single query. For example, `app.bsky.graph.listitem` records have both a `list` field and a `subject` field; joining them resolves list membership directly instead of requiring a `getBacklinks` call followed by N individual record lookups. + +| Parameter | Description | +|---|---| +| `subject` | Target AT-URI, DID, or URL (required) | +| `source` | Collection and field path (required) | +| `pathToOther` | Second field path on the same source record (required) | +| `linkDid` | Filter to specific linking-record DIDs (repeatable) | +| `otherSubject` | Filter to specific secondary link targets (repeatable) | +| `limit` | Page size, 1–1000 (default 100) | +| `cursor` | Pagination cursor from previous response | + +Response: `{"total": 42, "items": [{"linkRecord": {"did": "...", "collection": "...", "rkey": "..."}, "otherSubject": "..."}], "cursor": "..."}` + +### `GET /xrpc/blue.microcosm.links.getManyToManyCounts` + +Like `getManyToMany`, but grouped: counts of linking records per distinct secondary target instead of the individual records themselves. Useful when you only need aggregate counts, e.g. "how many people on each of these lists also follow me" without paginating every membership record. + +| Parameter | Description | +|---|---| +| `subject` | Target AT-URI, DID, or URL (required) | +| `source` | Collection and field path (required) | +| `pathToOther` | Second field path on the same source record (required) | +| `did` | Filter to specific linking-record DIDs (repeatable) | +| `otherSubject` | Filter to specific secondary link targets (repeatable) | +| `limit` | Page size, 1–1000 (default 100) | +| `cursor` | Pagination cursor from previous response | + +Response: `{"counts_by_other_subject": [{"subject": "...", "total": 42, "distinct": 12}], "cursor": "..."}` + +Note the DID filter parameter is `did` here, not `linkDid` like `getManyToMany` — a real inconsistency in the upstream Constellation API (their own source flags it as a known TODO), preserved here for compatibility rather than "fixed." + ## Roadmap **Near term** -- [ ] `blue.microcosm.links.getManyToMany` endpoint (Constellation parity) +- [x] Full Constellation API parity (`getBacklinksCount`, `getBacklinkDids`, `getBacklinks`, `getManyToMany`, `getManyToManyCounts`) - [ ] Configurable listen address, database path, and relay URL - [ ] Account deletion and deactivation handling - [x] Graceful shutdown and Firehose reconnect diff --git a/internal/api/many_to_many.go b/internal/api/many_to_many.go new file mode 100644 index 0000000..83d5346 --- /dev/null +++ b/internal/api/many_to_many.go @@ -0,0 +1,79 @@ +package api + +import ( + "encoding/base64" + "encoding/json" + "log" + "net/http" + "strconv" + + "github.com/alyraffauf/asterism/internal/store" +) + +type manyToManyResponse struct { + Total uint64 `json:"total"` + Items []store.ManyToManyItem `json:"items"` + Cursor *string `json:"cursor"` +} + +func (s *Server) GetManyToMany(w http.ResponseWriter, r *http.Request) { + query := r.URL.Query() + + subject := query.Get("subject") + source := query.Get("source") + linkDids := query["linkDid"] + otherSubjects := query["otherSubject"] + + collection, path, err := parseSource(source) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + pathToOther, err := normalizePath(query.Get("pathToOther")) + if err != nil { + http.Error(w, "path_to_other: "+err.Error(), http.StatusBadRequest) + return + } + + limit := uint64(100) + if raw := query.Get("limit"); raw != "" { + parsed, err := strconv.ParseUint(raw, 10, 64) + if err != nil || parsed == 0 || parsed > 1000 { + http.Error(w, "limit must be a number between 1 and 1000", http.StatusBadRequest) + return + } + limit = parsed + } + + var after int64 + if raw := query.Get("cursor"); raw != "" { + decoded, err := base64.StdEncoding.DecodeString(raw) + if err != nil { + http.Error(w, "invalid cursor", http.StatusBadRequest) + return + } + after, err = strconv.ParseInt(string(decoded), 10, 64) + if err != nil { + http.Error(w, "invalid cursor", http.StatusBadRequest) + return + } + } + + total, items, err := s.Store.ManyToMany(r.Context(), subject, collection, path, pathToOther, linkDids, otherSubjects, after, limit) + if err != nil { + log.Println("many to many:", err) + http.Error(w, "internal error", http.StatusInternalServerError) + return + } + + var cursor *string + if uint64(len(items)) == limit { + last := items[len(items)-1] + encoded := base64.StdEncoding.EncodeToString([]byte(strconv.FormatInt(last.LinkRecord.ID, 10))) + cursor = &encoded + } + + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(manyToManyResponse{Total: total, Items: items, Cursor: cursor}) +} diff --git a/internal/api/many_to_many_counts.go b/internal/api/many_to_many_counts.go new file mode 100644 index 0000000..649a115 --- /dev/null +++ b/internal/api/many_to_many_counts.go @@ -0,0 +1,73 @@ +package api + +import ( + "encoding/base64" + "encoding/json" + "log" + "net/http" + "strconv" + + "github.com/alyraffauf/asterism/internal/store" +) + +type manyToManyCountsResponse struct { + CountsByOtherSubject []store.OtherSubjectCount `json:"counts_by_other_subject"` + Cursor *string `json:"cursor"` +} + +func (s *Server) GetManyToManyCounts(w http.ResponseWriter, r *http.Request) { + query := r.URL.Query() + + subject := query.Get("subject") + source := query.Get("source") + linkDids := query["did"] + otherSubjects := query["otherSubject"] + + collection, path, err := parseSource(source) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + pathToOther, err := normalizePath(query.Get("pathToOther")) + if err != nil { + http.Error(w, "pathToOther: "+err.Error(), http.StatusBadRequest) + return + } + + limit := uint64(100) + if raw := query.Get("limit"); raw != "" { + parsed, err := strconv.ParseUint(raw, 10, 64) + if err != nil || parsed == 0 || parsed > 1000 { + http.Error(w, "limit must be a number between 1 and 1000", http.StatusBadRequest) + return + } + limit = parsed + } + + after := "" + if raw := query.Get("cursor"); raw != "" { + decoded, err := base64.StdEncoding.DecodeString(raw) + if err != nil { + http.Error(w, "invalid cursor", http.StatusBadRequest) + return + } + after = string(decoded) + } + + counts, err := s.Store.ManyToManyCounts(r.Context(), subject, collection, path, pathToOther, linkDids, otherSubjects, after, limit) + if err != nil { + log.Println("many to many counts:", err) + http.Error(w, "internal error", http.StatusInternalServerError) + return + } + + var cursor *string + if uint64(len(counts)) == limit { + encoded := base64.StdEncoding.EncodeToString([]byte(counts[len(counts)-1].Subject)) + cursor = &encoded + } + + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(manyToManyCountsResponse{CountsByOtherSubject: counts, Cursor: cursor}) +} diff --git a/internal/api/server.go b/internal/api/server.go index afb61b6..1bf88be 100644 --- a/internal/api/server.go +++ b/internal/api/server.go @@ -15,6 +15,8 @@ func (s *Server) Run(addr string) error { mux.HandleFunc("GET /xrpc/blue.microcosm.links.getBacklinksCount", s.GetBacklinksCount) mux.HandleFunc("GET /xrpc/blue.microcosm.links.getBacklinkDids", s.GetBacklinkDids) mux.HandleFunc("GET /xrpc/blue.microcosm.links.getBacklinks", s.GetBacklinks) + mux.HandleFunc("GET /xrpc/blue.microcosm.links.getManyToMany", s.GetManyToMany) + mux.HandleFunc("GET /xrpc/blue.microcosm.links.getManyToManyCounts", s.GetManyToManyCounts) return http.ListenAndServe(addr, mux) } diff --git a/internal/api/source.go b/internal/api/source.go index 17f437a..f67ce07 100644 --- a/internal/api/source.go +++ b/internal/api/source.go @@ -14,14 +14,23 @@ func parseSource(raw string) (collection string, path string, err error) { return "", "", fmt.Errorf("source is missing a collection before ':'") } + path, err = normalizePath(rawPath) + if err != nil { + return "", "", fmt.Errorf("source path: %w", err) + } + + return collection, path, nil +} + +func normalizePath(rawPath string) (string, error) { switch { case rawPath == "": - return "", "", fmt.Errorf("source is missing a path after ':'") + return "", fmt.Errorf("path is empty") case rawPath == ".": - return collection, ".", nil + return ".", nil case strings.HasPrefix(rawPath, "."): - return "", "", fmt.Errorf("source path must not start with '.'") + return "", fmt.Errorf("path must not start with '.'") default: - return collection, "." + rawPath, nil + return "." + rawPath, nil } } diff --git a/internal/store/query.go b/internal/store/query.go index 90ae5d2..f99fda9 100644 --- a/internal/store/query.go +++ b/internal/store/query.go @@ -13,6 +13,18 @@ type Record struct { RecordKey string `json:"rkey"` } +type ManyToManyItem struct { + LinkRecord Record `json:"linkRecord"` + OtherSubject string `json:"otherSubject"` +} + + +type OtherSubjectCount struct { + Subject string `json:"subject"` + Total uint64 `json:"total"` + Distinct uint64 `json:"distinct"` +} + func (s *Store) CountBacklinks(ctx context.Context, target, collection, fieldPath string) (uint64, error) { var total uint64 @@ -122,3 +134,118 @@ func (s *Store) ListBacklinks(ctx context.Context, target, collection, fieldPath return total, records, nil } + +func (s *Store) ManyToMany(ctx context.Context, target, collection, fieldPath, pathToOther string, linkDids, otherSubjects []string, after int64, limit uint64) (total uint64, items []ManyToManyItem, err error) { + where := `a.target = ? AND a.collection = ? AND a.field_path = ? AND b.field_path = ?` + args := []any{target, collection, fieldPath, pathToOther} + + if len(linkDids) > 0 { + placeholders := make([]string, len(linkDids)) + for i, did := range linkDids { + placeholders[i] = "?" + args = append(args, did) + } + where += ` AND a.actor_did IN (` + strings.Join(placeholders, ", ") + `)` + } + + if len(otherSubjects) > 0 { + placeholders := make([]string, len(otherSubjects)) + for i, subj := range otherSubjects { + placeholders[i] = "?" + args = append(args, subj) + } + where += ` AND b.target IN (` + strings.Join(placeholders, ", ") + `)` + } + + joinClause := ` FROM links a JOIN links b + ON a.actor_did = b.actor_did AND a.collection = b.collection AND a.record_key = b.record_key + WHERE ` + where + + if err := s.readDB.QueryRowContext(ctx, `SELECT COUNT(*)`+joinClause, args...).Scan(&total); err != nil { + return 0, nil, fmt.Errorf("count many to many: %w", err) + } + + listQuery := `SELECT a.id, a.actor_did, a.collection, a.record_key, b.target` + joinClause + listArgs := append([]any{}, args...) + + if after != 0 { + listQuery += ` AND a.id > ?` + listArgs = append(listArgs, after) + } + + listQuery += ` ORDER BY a.id ASC LIMIT ?` + listArgs = append(listArgs, limit) + + rows, err := s.readDB.QueryContext(ctx, listQuery, listArgs...) + if err != nil { + return 0, nil, fmt.Errorf("query many to many: %w", err) + } + defer rows.Close() + + for rows.Next() { + var item ManyToManyItem + if err := rows.Scan(&item.LinkRecord.ID, &item.LinkRecord.ActorDid, &item.LinkRecord.Collection, &item.LinkRecord.RecordKey, &item.OtherSubject); err != nil { + return 0, nil, fmt.Errorf("scan many to many item: %w", err) + } + items = append(items, item) + } + if err := rows.Err(); err != nil { + return 0, nil, fmt.Errorf("iterate many to many: %w", err) + } + + return total, items, nil +} + +func (s *Store) ManyToManyCounts(ctx context.Context, target, collection, fieldPath, pathToOther string, linkDids, otherSubjects []string, after string, limit uint64) (counts []OtherSubjectCount, err error) { + where := `a.target = ? AND a.collection = ? AND a.field_path = ? AND b.field_path = ?` + args := []any{target, collection, fieldPath, pathToOther} + + if len(linkDids) > 0 { + placeholders := make([]string, len(linkDids)) + for i, did := range linkDids { + placeholders[i] = "?" + args = append(args, did) + } + where += ` AND a.actor_did IN (` + strings.Join(placeholders, ", ") + `)` + } + + if len(otherSubjects) > 0 { + placeholders := make([]string, len(otherSubjects)) + for i, subj := range otherSubjects { + placeholders[i] = "?" + args = append(args, subj) + } + where += ` AND b.target IN (` + strings.Join(placeholders, ", ") + `)` + } + + where += ` AND b.target > ?` + args = append(args, after) + + query := `SELECT b.target, COUNT(*), COUNT(DISTINCT a.actor_did) + FROM links a JOIN links b + ON a.actor_did = b.actor_did AND a.collection = b.collection AND a.record_key = b.record_key + WHERE ` + where + ` + GROUP BY b.target + ORDER BY b.target + LIMIT ?` + args = append(args, limit) + + rows, err := s.readDB.QueryContext(ctx, query, args...) + if err != nil { + return nil, fmt.Errorf("query many to many counts: %w", err) + } + defer rows.Close() + + for rows.Next() { + var c OtherSubjectCount + if err := rows.Scan(&c.Subject, &c.Total, &c.Distinct); err != nil { + return nil, fmt.Errorf("scan many to many count: %w", err) + } + counts = append(counts, c) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterate many to many counts: %w", err) + } + + return counts, nil +} -- 2.51.2