From bc5c54c72f5affcb138e3c796da6dbaac08b9b5f Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Sat, 25 Apr 2026 05:15:27 +0900 Subject: [PATCH] appview: use blob based strings Signed-off-by: Seongmin Lee --- appview/db/db.go | 18 +++ appview/db/strings.go | 156 ++++++++++++++------ appview/ingester.go | 43 +++++- appview/ingester_string_test.go | 20 +-- appview/models/string.go | 70 ++++++--- appview/state/router.go | 13 +- appview/state/state.go | 6 + appview/strings/strings.go | 247 +++++++++++++++++++++++--------- blobstore/blobstore.go | 13 ++ blobstore/pds.go | 46 ++++++ 10 files changed, 484 insertions(+), 148 deletions(-) create mode 100644 blobstore/blobstore.go create mode 100644 blobstore/pds.go diff --git a/appview/db/db.go b/appview/db/db.go index 052c421f..6528a719 100644 --- a/appview/db/db.go +++ b/appview/db/db.go @@ -2152,6 +2152,24 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { return nil }) + orm.RunMigration(conn, logger, "use-blobs-for-string-files", func(tx *sql.Tx) error { + _, err := tx.Exec(` + create table string_files ( + id integer primary key autoincrement, + at_uri text not null, + name text not null, + content_ref text not null, + content_size integer not null, + content_mimetype text not null, + gzip_realsize integer, + gzip_realmime text, + gzip_realcontent text, -- decoded content of gzipped text blobs + foreign key (at_uri) references strings(at_uri) on delete cascade + ); + `) + return err + }) + return &DB{ db, logger, diff --git a/appview/db/strings.go b/appview/db/strings.go index 2a4a4baa..e7d8d8a9 100644 --- a/appview/db/strings.go +++ b/appview/db/strings.go @@ -9,6 +9,9 @@ import ( "time" "github.com/bluesky-social/indigo/atproto/syntax" + lexutil "github.com/bluesky-social/indigo/lex/util" + "github.com/ipfs/go-cid" + "tangled.org/core/api/tangled" "tangled.org/core/appview/models" "tangled.org/core/orm" ) @@ -62,6 +65,55 @@ func AddString(d *DB, s models.String) error { return nil } + _, err = tx.Exec(`delete from string_files where at_uri = ?`, s.AtUri()) + if err != nil { + return fmt.Errorf("deleting old files: %w", err) + } + + vals := make([]string, len(s.Files)) + args := make([]any, 0, len(s.Files)*8) + for i, file := range s.Files { + vals[i] = "(?, ?, ?, ?, ?, ?, ?, ?)" + var gzipRealSize *int64 + var gzipRealMime *string + var gzipRealContent *string + if file.Gzip != nil { + gzipRealSize = &file.Gzip.RealSize + gzipRealMime = &file.Gzip.RealMime + gzipRealContent = &file.Gzip.Content + } + args = append(args, + s.AtUri(), + file.Name, + file.Content.Ref.String(), + file.Content.Size, + file.Content.MimeType, + gzipRealSize, + gzipRealMime, + gzipRealContent, + ) + } + _, err = tx.Exec( + fmt.Sprintf( + `insert into string_files ( + at_uri, + name, + content_ref, + content_size, + content_mimetype, + gzip_realsize, + gzip_realmime, + gzip_realcontent + ) + values %s`, + strings.Join(vals, ","), + ), + args..., + ) + if err != nil { + return fmt.Errorf("inserting files: %w", err) + } + if err := tx.Commit(); err != nil { return fmt.Errorf("commiting transaction: %w", err) } @@ -194,47 +246,69 @@ func GetStrings(e Execer, limit int, filters ...orm.Filter) ([]models.String, er i++ } - // // get files - // { - // rows, err := e.Query( - // fmt.Sprintf( - // `select at_uri, name, blob from string_files where at_uri in (%s) order by at_uri, id`, - // inClause, - // ), - // args..., - // ) - // if err != nil { - // return nil, fmt.Errorf("failed to execute string_files query: %w", err) - // } - // defer rows.Close() - // - // for rows.Next() { - // var stringAt syntax.ATURI - // var file models.String_File - // file.Blob = &util.LexBlob{} - // var gzipMimeType sql.Null[string] - // var gzipSize sql.Null[int64] - // if err := rows.Scan( - // &stringAt, - // &file.Name, - // &blob, - // ); err != nil { - // return nil, fmt.Errorf("failed to execute string_files query: %w", err) - // } - // if gzipMimeType.Valid && gzipSize.Valid { - // file.Gzip = &models.GzipInfo{ - // MimeType: gzipMimeType.V, - // Size: gzipSize.V, - // } - // } - // if s, ok := stringMap[stringAt]; ok { - // s.Files = append(s.Files, file) - // } - // } - // if err = rows.Err(); err != nil { - // return nil, fmt.Errorf("failed to execute string_files query: %w", err) - // } - // } + // get files + { + rows, err := e.Query( + fmt.Sprintf( + `select + at_uri, + name, + content_ref, + content_size, + content_mimetype, + gzip_realsize, + gzip_realmime, + gzip_realcontent + from string_files + where at_uri in (%s) order by at_uri, id`, + inClause, + ), + args..., + ) + if err != nil { + return nil, fmt.Errorf("failed to execute string_files query: %w", err) + } + defer rows.Close() + + for rows.Next() { + var stringAt syntax.ATURI + var file models.String_File + + var contentRef string + var gzipRealSize sql.Null[int64] + var gzipRealMime, gzipRealContent sql.Null[string] + if err := rows.Scan( + &stringAt, + &file.Name, + &contentRef, + &file.Content.Size, + &file.Content.MimeType, + &gzipRealSize, + &gzipRealMime, + &gzipRealContent, + ); err != nil { + return nil, fmt.Errorf("failed to execute string_files query: %w", err) + } + + file.Content.Ref = lexutil.LexLink(cid.MustParse(contentRef)) + + if gzipRealMime.Valid && gzipRealSize.Valid && gzipRealContent.Valid { + file.Gzip = &models.String_GzipInfo{ + String_File_Gzip: tangled.String_File_Gzip{ + RealMime: gzipRealMime.V, + RealSize: gzipRealSize.V, + }, + Content: gzipRealContent.V, + } + } + if s, ok := stringMap[stringAt]; ok { + s.Files = append(s.Files, file) + } + } + if err = rows.Err(); err != nil { + return nil, fmt.Errorf("failed to execute string_files query: %w", err) + } + } // get star counts { diff --git a/appview/ingester.go b/appview/ingester.go index 74fcec90..685147f4 100644 --- a/appview/ingester.go +++ b/appview/ingester.go @@ -1,6 +1,7 @@ package appview import ( + "compress/gzip" "context" "database/sql" "encoding/json" @@ -33,6 +34,7 @@ import ( "tangled.org/core/appview/repoverify" "tangled.org/core/appview/serververify" "tangled.org/core/appview/validator" + "tangled.org/core/blobstore" "tangled.org/core/idresolver" "tangled.org/core/orm" "tangled.org/core/rbac" @@ -43,6 +45,7 @@ type Ingester struct { Db *db.DB Enforcer *rbac.Enforcer IdResolver *idresolver.Resolver + BlobStore blobstore.BlobStore Cache *cache.Cache Config *config.Config Logger *slog.Logger @@ -94,7 +97,7 @@ func (i *Ingester) Ingest() processFunc { case tangled.KnotNSID: err = i.ingestKnot(ctx, e) case tangled.StringNSID: - err = i.ingestString(e) + err = i.ingestString(ctx, e) case tangled.RepoIssueNSID: err = i.ingestIssue(ctx, e) case tangled.RepoPullNSID: @@ -903,7 +906,7 @@ func (i *Ingester) ingestSpindle(ctx context.Context, e *jmodels.Event) error { return nil } -func (i *Ingester) ingestString(e *jmodels.Event) error { +func (i *Ingester) ingestString(ctx context.Context, e *jmodels.Event) error { did := e.Did rkey := e.Commit.RKey @@ -922,16 +925,46 @@ func (i *Ingester) ingestString(e *jmodels.Event) error { return err } - string, err := models.StringFromRecord(syntax.DID(did), syntax.RecordKey(rkey), syntax.CID(e.Commit.CID), record) + str, err := models.StringFromRecord(syntax.DID(did), syntax.RecordKey(rkey), syntax.CID(e.Commit.CID), record) if err != nil { return fmt.Errorf("failed to parse string record: %w", err) } - if err = string.Validate(); err != nil { + if err = str.Validate(); err != nil { l.Error("invalid record", "err", err) return err } - if err = db.AddString(i.Db, string); err != nil { + g, gctx := errgroup.WithContext(ctx) + for idx, file := range str.Files { + if file.Gzip == nil { + continue + } + g.Go(func() error { + blob, err := i.BlobStore.GetBlob(gctx, str.Did, cid.Cid(file.Content.Ref)) + if err != nil { + return fmt.Errorf("files[%d]: failed to fetch blob: %w", idx, err) + } + defer blob.Close() + gzr, err := gzip.NewReader(blob) + if err != nil { + return fmt.Errorf("files[%d]: invalid gzip stream: %w", idx, err) + } + gzr.Close() + + content, err := io.ReadAll(gzr) + if err != nil { + return fmt.Errorf("files[%d]: failed to read blob: %w", idx, err) + } + file.Gzip.Content = string(content) + str.Files[idx] = file + return nil + }) + } + if err := g.Wait(); err != nil { + return err + } + + if err = db.AddString(i.Db, str); err != nil { l.Error("failed to add string", "err", err) return err } diff --git a/appview/ingester_string_test.go b/appview/ingester_string_test.go index ecde8bf1..897c3956 100644 --- a/appview/ingester_string_test.go +++ b/appview/ingester_string_test.go @@ -96,7 +96,7 @@ func TestIngestString_CreateRoundTrip(t *testing.T) { CreatedAt: created.Format(time.RFC3339), }) - if err := ing.ingestString(e); err != nil { + if err := ing.ingestString(t.Context(), e); err != nil { t.Fatalf("ingestString: %v", err) } @@ -134,13 +134,13 @@ func TestIngestString_UpdateBumpsEditedOnContentChange(t *testing.T) { CreatedAt: created.Format(time.RFC3339), } - if err := ing.ingestString(makeStringEvent(t, jmodels.CommitOperationCreate, "did:plc:boltless", "rk1", base)); err != nil { + if err := ing.ingestString(t.Context(), makeStringEvent(t, jmodels.CommitOperationCreate, "did:plc:boltless", "rk1", base)); err != nil { t.Fatalf("ingestString create: %v", err) } updated := base updated.Contents = "hello, world!\n" - if err := ing.ingestString(makeStringEvent(t, jmodels.CommitOperationUpdate, "did:plc:boltless", "rk1", updated)); err != nil { + if err := ing.ingestString(t.Context(), makeStringEvent(t, jmodels.CommitOperationUpdate, "did:plc:boltless", "rk1", updated)); err != nil { t.Fatalf("ingestString update: %v", err) } @@ -168,10 +168,10 @@ func TestIngestString_UpdateNoChangeKeepsEditedNil(t *testing.T) { CreatedAt: time.Date(2025, 9, 14, 10, 30, 0, 0, time.UTC).Format(time.RFC3339), } - if err := ing.ingestString(makeStringEvent(t, jmodels.CommitOperationCreate, "did:plc:akshay", "rk2", rec)); err != nil { + if err := ing.ingestString(t.Context(), makeStringEvent(t, jmodels.CommitOperationCreate, "did:plc:akshay", "rk2", rec)); err != nil { t.Fatalf("create: %v", err) } - if err := ing.ingestString(makeStringEvent(t, jmodels.CommitOperationUpdate, "did:plc:akshay", "rk2", rec)); err != nil { + if err := ing.ingestString(t.Context(), makeStringEvent(t, jmodels.CommitOperationUpdate, "did:plc:akshay", "rk2", rec)); err != nil { t.Fatalf("update: %v", err) } @@ -188,7 +188,7 @@ func TestIngestString_DeleteRemovesRow(t *testing.T) { Contents: "x", CreatedAt: time.Now().UTC().Format(time.RFC3339), } - if err := ing.ingestString(makeStringEvent(t, jmodels.CommitOperationCreate, "did:plc:boltless", "rk1", rec)); err != nil { + if err := ing.ingestString(t.Context(), makeStringEvent(t, jmodels.CommitOperationCreate, "did:plc:boltless", "rk1", rec)); err != nil { t.Fatalf("create: %v", err) } @@ -202,7 +202,7 @@ func TestIngestString_DeleteRemovesRow(t *testing.T) { CID: makeCID(t, &rec), }, } - if err := ing.ingestString(del); err != nil { + if err := ing.ingestString(t.Context(), del); err != nil { t.Fatalf("delete: %v", err) } @@ -226,7 +226,7 @@ func TestIngestString_ValidatorRejects(t *testing.T) { for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { e := makeStringEvent(t, jmodels.CommitOperationCreate, "did:plc:akshay", "bad", tc.rec) - if err := ing.ingestString(e); err == nil { + if err := ing.ingestString(t.Context(), e); err == nil { t.Fatal("expected validator error, got nil") } if _, ok := loadString(t, ing, "did:plc:akshay", "bad"); ok { @@ -247,13 +247,13 @@ func TestIngestString_ColdReplayPreservesCreated(t *testing.T) { event := makeStringEvent(t, jmodels.CommitOperationCreate, "did:plc:boltless", "rkcold", rec) first := newStringIngester(t) - if err := first.ingestString(event); err != nil { + if err := first.ingestString(t.Context(), event); err != nil { t.Fatalf("first ingest: %v", err) } live, _ := loadString(t, first, "did:plc:boltless", "rkcold") second := newStringIngester(t) - if err := second.ingestString(event); err != nil { + if err := second.ingestString(t.Context(), event); err != nil { t.Fatalf("replay ingest: %v", err) } replayed, _ := loadString(t, second, "did:plc:boltless", "rkcold") diff --git a/appview/models/string.go b/appview/models/string.go index f377fece..9711944a 100644 --- a/appview/models/string.go +++ b/appview/models/string.go @@ -37,15 +37,17 @@ type String struct { Stats *StringStats } -// TODO: replace this with [tangled.String_File] +// String_File is [tangled.String_File] with optional prefetched & decompressed text blob content type String_File struct { - Name string - Blob *lexutil.LexBlob - Gzip *GzipInfo + Name string + Content lexutil.LexBlob + Gzip *String_GzipInfo } -type GzipInfo struct { - MimeType string - Size int64 +type String_GzipInfo struct { + tangled.String_File_Gzip + // Optional uncompressed content. + // Populated when the content is first requested. + Content string } func (s *String) AtUri() syntax.ATURI { @@ -53,14 +55,25 @@ func (s *String) AtUri() syntax.ATURI { } func (s *String) AsRecord() *tangled.String { - var description string - if s.Description != nil { - description = *s.Description + var files []*tangled.String_File + for _, f := range s.Files { + var gzip *tangled.String_File_Gzip + if f.Gzip != nil { + gzip = &tangled.String_File_Gzip{ + RealSize: f.Gzip.RealSize, + RealMime: f.Gzip.RealMime, + } + } + files = append(files, &tangled.String_File{ + Name: f.Name, + Content: &f.Content, + Gzip: gzip, + }) } return &tangled.String{ - Filename: s.FileName, - Description: description, - Contents: s.FileContent, + Title: s.Title, + Description: s.Description, + Files: files, CreatedAt: s.Created.Format(time.RFC3339), } } @@ -116,23 +129,35 @@ func (s String) IsLegacySingleFile() bool { return len(s.Files) == 0 } +// StringFromRecord creates [String] from [tangled.String]. +// NOTE: This won't prefetch blobs func StringFromRecord(did syntax.DID, rkey syntax.RecordKey, cid syntax.CID, record tangled.String) (String, error) { created, err := time.Parse(time.RFC3339, record.CreatedAt) if err != nil { return String{}, fmt.Errorf("invalid createdAt: %w", err) } - var description *string - if record.Description != "" { - description = &record.Description + var files []String_File + for _, f := range record.Files { + var gzip *String_GzipInfo + if f.Gzip != nil { + gzip = &String_GzipInfo{String_File_Gzip: *f.Gzip} + } + files = append(files, String_File{ + Name: f.Name, + Content: *f.Content, + Gzip: gzip, + }) } return String{ Did: did, Rkey: rkey, Cid: &cid, - Description: description, + Title: record.Title, + Description: record.Description, + Files: files, Created: created, - FileName: record.Filename, - FileContent: record.Contents, + FileName: stringPtr(record.Filename), + FileContent: stringPtr(record.Contents), }, nil } @@ -145,3 +170,10 @@ type StringFileStats struct { LineCount int ByteCount int } + +func stringPtr(s *string) string { + if s == nil { + return "" + } + return *s +} diff --git a/appview/state/router.go b/appview/state/router.go index f06776a9..d8ed3212 100644 --- a/appview/state/router.go +++ b/appview/state/router.go @@ -323,12 +323,13 @@ func (s *State) StringsRouter(mw *middleware.Middleware) http.Handler { logger := log.SubLogger(s.logger, "strings") strs := &avstrings.Strings{ - Db: s.db, - OAuth: s.oauth, - Pages: s.pages, - Dir: s.idResolver.Directory(), - Notifier: s.notifier, - Logger: logger, + Db: s.db, + OAuth: s.oauth, + Pages: s.pages, + Dir: s.idResolver.Directory(), + Notifier: s.notifier, + Logger: logger, + BlobStore: s.blobStore, } return strs.Router(mw) diff --git a/appview/state/state.go b/appview/state/state.go index 21659236..dec95262 100644 --- a/appview/state/state.go +++ b/appview/state/state.go @@ -32,6 +32,7 @@ import ( "tangled.org/core/appview/repoverify" "tangled.org/core/appview/validator" xrpcclient "tangled.org/core/appview/xrpcclient" + "tangled.org/core/blobstore" "tangled.org/core/consts" "tangled.org/core/eventconsumer" "tangled.org/core/idresolver" @@ -59,6 +60,7 @@ type State struct { enforcer *rbac.Enforcer pages *pages.Pages idResolver *idresolver.Resolver + blobStore blobstore.BlobStore rdb *cache.Cache mentionsResolver *mentions.Resolver posthog posthog.Client @@ -118,6 +120,8 @@ func Make(ctx context.Context, config *config.Config) (*State, error) { mentionsResolver := mentions.New(config, res, d, log.SubLogger(logger, "mentionsResolver")) + blobStore := blobstore.NewPdsBlobStore(res.Directory()) + jc, err := jetstream.NewJetstreamClient( config.Jetstream.Endpoint, "appview", @@ -180,6 +184,7 @@ func Make(ctx context.Context, config *config.Config) (*State, error) { Db: d, Enforcer: enforcer, IdResolver: res, + BlobStore: blobStore, Cache: rdb, Config: config, Logger: log.SubLogger(logger, "ingester"), @@ -224,6 +229,7 @@ func Make(ctx context.Context, config *config.Config) (*State, error) { enforcer: enforcer, pages: pages, idResolver: res, + blobStore: blobStore, rdb: rdb, mentionsResolver: mentionsResolver, posthog: posthog, diff --git a/appview/strings/strings.go b/appview/strings/strings.go index a6a7d416..7ae30f5c 100644 --- a/appview/strings/strings.go +++ b/appview/strings/strings.go @@ -2,6 +2,7 @@ package stringn import ( "bytes" + "compress/gzip" "context" "database/sql" "errors" @@ -21,27 +22,33 @@ import ( "tangled.org/core/appview/oauth" "tangled.org/core/appview/pages" "tangled.org/core/appview/pages/markup" + "tangled.org/core/blobstore" "tangled.org/core/orm" "tangled.org/core/tid" + "tangled.org/core/xrpc" "github.com/bluesky-social/indigo/api/agnostic" "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" - "github.com/bluesky-social/indigo/xrpc" + indigoxrpc "github.com/bluesky-social/indigo/xrpc" "github.com/go-chi/chi/v5" + "github.com/ipfs/go-cid" comatproto "github.com/bluesky-social/indigo/api/atproto" lexutil "github.com/bluesky-social/indigo/lex/util" ) +const ApplicationGzip = "application/gzip" + type Strings struct { - Db *db.DB - OAuth *oauth.OAuth - Pages *pages.Pages - Dir identity.Directory - Logger *slog.Logger - Notifier notify.Notifier + Db *db.DB + OAuth *oauth.OAuth + Pages *pages.Pages + Dir identity.Directory + BlobStore blobstore.BlobStore + Logger *slog.Logger + Notifier notify.Notifier } func (s *Strings) Router(mw *middleware.Middleware) http.Handler { @@ -117,7 +124,7 @@ func (s *Strings) resolveString(next http.Handler) http.Handler { return } - string, err := db.GetString(s.Db, orm.FilterEq("did", id.DID), orm.FilterEq("rkey", rkey)) + str, err := db.GetString(s.Db, orm.FilterEq("did", id.DID), orm.FilterEq("rkey", rkey)) if errors.Is(err, sql.ErrNoRows) { s.Pages.Error404(w) return @@ -127,7 +134,7 @@ func (s *Strings) resolveString(next http.Handler) http.Handler { return } - ctx := context.WithValue(r.Context(), stringCtxKey{}, string) + ctx := context.WithValue(r.Context(), stringCtxKey{}, str) next.ServeHTTP(w, r.WithContext(ctx)) }) } @@ -170,24 +177,24 @@ func (s *Strings) SingleString(w http.ResponseWriter, r *http.Request) { l := s.Logger.With("handler", "SingleString") ctx := r.Context() - string, ok := stringFromContext(ctx) + str, ok := stringFromContext(ctx) if !ok { l.Error("malformed middleware. string missing") s.Pages.Error404(w) return } - starCount, err := db.GetStarCount(s.Db, models.StarSubjectString, string.AtUri().String()) + starCount, err := db.GetStarCount(s.Db, models.StarSubjectString, str.AtUri().String()) if err != nil { l.Error("failed to get star count", "err", err) } user := s.OAuth.GetMultiAccountUser(r) isStarred := false if user != nil { - isStarred = db.GetStarStatus(s.Db, user.Did, string.AtUri().String()) + isStarred = db.GetStarStatus(s.Db, user.Did, str.AtUri().String()) } - comments, err := db.GetComments(s.Db, orm.FilterEq("subject_uri", string.AtUri())) + comments, err := db.GetComments(s.Db, orm.FilterEq("subject_uri", str.AtUri())) if err != nil { l.Error("failed to get comments", "err", err) } @@ -223,22 +230,39 @@ func (s *Strings) SingleString(w http.ResponseWriter, r *http.Request) { var files []pages.StringFileFragmentParams - if string.IsLegacySingleFile() { + if str.IsLegacySingleFile() { files = []pages.StringFileFragmentParams{ - s.makeFileFragmentParams(&string, string.FileName, string.FileContent, false), + s.makeFileFragmentParams(&str, str.FileName, str.FileContent, false), } } else { - files = make([]pages.StringFileFragmentParams, len(string.Files)) - for i, file := range string.Files { - // TODO: read blob - content := "" - files[i] = s.makeFileFragmentParams(&string, file.Name, content, false) + files = make([]pages.StringFileFragmentParams, len(str.Files)) + for i, file := range str.Files { + var content string + if file.Gzip != nil { + content = file.Gzip.Content + } else { + blob, err := s.BlobStore.GetBlob(r.Context(), str.Did, cid.Cid(file.Content.Ref)) + if err != nil { + l.Warn("failed to fetch blob", "err", err) + http.NotFound(w, r) + return + } + defer blob.Close() + + contentBytes, err := io.ReadAll(blob) + if err != nil { + l.Error("failed to read blob", "err", err) + } + content = string(contentBytes) + } + + files[i] = s.makeFileFragmentParams(&str, file.Name, content, false) } } err = s.Pages.SingleString(w, pages.SingleStringParams{ LoggedInUser: user, - String: &string, + String: &str, FileParams: files, IsStarred: isStarred, StarCount: starCount, @@ -288,12 +312,28 @@ func (s *Strings) edit(w http.ResponseWriter, r *http.Request) { } else { files = make([]pages.StringFileEditFragmentParams, len(oldString.Files)) for i, file := range oldString.Files { - // TODO: read blob - content := "" + var content string + if file.Gzip != nil { + content = file.Gzip.Content + } else { + blob, err := s.BlobStore.GetBlob(r.Context(), oldString.Did, cid.Cid(file.Content.Ref)) + if err != nil { + l.Warn("failed to fetch blob", "err", err) + http.NotFound(w, r) + return + } + defer blob.Close() + + contentBytes, err := io.ReadAll(blob) + if err != nil { + l.Error("failed to read blob", "err", err) + } + content = string(contentBytes) + } files[i] = pages.StringFileEditFragmentParams{ Name: file.Name, Content: content, - Size: uint64(file.Blob.Size), + Size: uint64(file.Content.Size), } } } @@ -333,18 +373,35 @@ func (s *Strings) edit(w http.ResponseWriter, r *http.Request) { return } - newString := oldString - newString.Title = title - newString.Description = description - newString.FileName = filename - newString.FileContent = content - client, err := s.OAuth.AuthorizedClient(r) if err != nil { fail("Failed to create record.", err) return } + blob, err := xrpc.RepoUploadBlob(ctx, client, gz(content), ApplicationGzip) + if err != nil { + fail("Failed to create record.", err) + return + } + + newString := oldString + newString.Title = title + newString.Description = description + newString.Files = []models.String_File{ + { + Name: filename, + Content: *blob.Blob, + Gzip: &models.String_GzipInfo{ + String_File_Gzip: tangled.String_File_Gzip{ + RealMime: "text/plain", + RealSize: int64(len(content)), + }, + Content: content, + }, + }, + } + // first replace the existing record in the PDS var exCid string if newString.Cid != nil { @@ -390,6 +447,7 @@ func (s *Strings) edit(w http.ResponseWriter, r *http.Request) { func (s *Strings) create(w http.ResponseWriter, r *http.Request) { l := s.Logger.With("handler", "create") + ctx := r.Context() user := s.OAuth.GetMultiAccountUser(r) switch r.Method { @@ -428,45 +486,62 @@ func (s *Strings) create(w http.ResponseWriter, r *http.Request) { return } - string := models.String{ - Did: syntax.DID(user.Did), - Rkey: syntax.RecordKey(tid.TID()), - Title: title, - Description: description, - FileName: filename, - FileContent: content, - Created: time.Now(), + client, err := s.OAuth.AuthorizedClient(r) + if err != nil { + fail("Failed to create record.", err) + return } - client, err := s.OAuth.AuthorizedClient(r) + blob, err := xrpc.RepoUploadBlob(ctx, client, gz(content), ApplicationGzip) if err != nil { fail("Failed to create record.", err) return } - resp, err := comatproto.RepoPutRecord(r.Context(), client, &atproto.RepoPutRecord_Input{ + newString := models.String{ + Did: syntax.DID(user.Did), + Rkey: syntax.RecordKey(tid.TID()), + Title: title, + Description: description, + Files: []models.String_File{ + { + Name: filename, + Content: *blob.Blob, + Gzip: &models.String_GzipInfo{ + String_File_Gzip: tangled.String_File_Gzip{ + RealMime: "text/plain", + RealSize: int64(len(content)), + }, + Content: content, + }, + }, + }, + Created: time.Now(), + } + + resp, err := comatproto.RepoPutRecord(ctx, client, &atproto.RepoPutRecord_Input{ Collection: tangled.StringNSID, - Repo: string.Did.String(), - Rkey: string.Rkey.String(), - Record: &lexutil.LexiconTypeDecoder{Val: string.AsRecord()}, + Repo: newString.Did.String(), + Rkey: newString.Rkey.String(), + Record: &lexutil.LexiconTypeDecoder{Val: newString.AsRecord()}, }) if err != nil { fail("Failed to create record.", err) return } l := l.With("aturi", resp.Uri) - l.Info("created record") + l.Info("created record", "files", len(newString.Files)) // insert into DB - if err = db.AddString(s.Db, string); err != nil { + if err = db.AddString(s.Db, newString); err != nil { fail("Failed to create string.", err) return } - s.Notifier.NewString(r.Context(), &string) + s.Notifier.NewString(ctx, &newString) // successful - s.Pages.HxRedirect(w, fmt.Sprintf("/strings/%s/%s", string.Did, string.Rkey)) + s.Pages.HxRedirect(w, fmt.Sprintf("/strings/%s/%s", newString.Did, newString.Rkey)) } } @@ -533,7 +608,7 @@ func (s *Strings) FileRaw(w http.ResponseWriter, r *http.Request) { l := s.Logger.With("handler", "FileRaw") ctx := r.Context() - string, ok := stringFromContext(ctx) + str, ok := stringFromContext(ctx) if !ok { l.Error("malformed middleware. string missing") s.Pages.Error404(w) @@ -541,32 +616,46 @@ func (s *Strings) FileRaw(w http.ResponseWriter, r *http.Request) { } filename := chi.URLParam(r, "filename") - if string.IsLegacySingleFile() { - if filename != string.FileName { + if str.IsLegacySingleFile() { + if filename != str.FileName { http.NotFound(w, r) return } w.Header().Set("Content-Type", "text/plain; charset=utf-8") w.Header().Set("Content-Disposition", fmt.Sprintf("inline; filename=%q", filename)) - w.Header().Set("Content-Length", strconv.Itoa(len(string.FileContent))) - _, err := w.Write([]byte(string.FileContent)) + w.Header().Set("Content-Length", strconv.Itoa(len(str.FileContent))) + _, err := w.Write([]byte(str.FileContent)) if err != nil { l.Error("failed to write raw response", "err", err) } } else { - file, ok := string.FileByName(filename) + file, ok := str.FileByName(filename) if !ok { http.NotFound(w, r) return } - content := "" + mimeType := file.Content.MimeType + size := file.Content.Size - w.Header().Set("Content-Type", "text/plain; charset=utf-8") + var reader io.Reader + if file.Gzip != nil { + reader = strings.NewReader(file.Gzip.Content) + } else { + blob, err := s.BlobStore.GetBlob(r.Context(), str.Did, cid.Cid(file.Content.Ref)) + if err != nil { + l.Warn("failed to fetch blob", "err", err) + http.NotFound(w, r) + return + } + defer blob.Close() + reader = blob + } + + w.Header().Set("Content-Type", mimeType) w.Header().Set("Content-Disposition", fmt.Sprintf("inline; filename=%q", filename)) - w.Header().Set("Content-Length", strconv.FormatInt(file.Blob.Size, 10)) - _, err := w.Write([]byte(content)) - if err != nil { + w.Header().Set("Content-Length", strconv.FormatInt(size, 10)) + if _, err := io.Copy(w, reader); err != nil { l.Error("failed to write raw response", "err", err) } } @@ -611,7 +700,7 @@ func (s *Strings) FileFragment(w http.ResponseWriter, r *http.Request) { l := s.Logger.With("handler", "FileFragment") ctx := r.Context() - string, ok := stringFromContext(ctx) + str, ok := stringFromContext(ctx) if !ok { l.Error("malformed middleware. string missing") http.NotFound(w, r) @@ -621,24 +710,40 @@ func (s *Strings) FileFragment(w http.ResponseWriter, r *http.Request) { forceCode := r.URL.Query().Get("code") == "true" var params pages.StringFileFragmentParams - if string.IsLegacySingleFile() { - if filename != string.FileName { + if str.IsLegacySingleFile() { + if filename != str.FileName { http.NotFound(w, r) return } - params = s.makeFileFragmentParams(&string, string.FileName, string.FileContent, forceCode) + params = s.makeFileFragmentParams(&str, str.FileName, str.FileContent, forceCode) } else { - file, ok := string.FileByName(filename) + file, ok := str.FileByName(filename) if !ok { l.Error("malformed middleware. string missing") http.NotFound(w, r) return } - // TODO: read blob - content := "" + var content string + if file.Gzip != nil && file.Gzip.Content != "" { + content = file.Gzip.Content + } else { + blob, err := s.BlobStore.GetBlob(r.Context(), str.Did, cid.Cid(file.Content.Ref)) + if err != nil { + l.Warn("failed to fetch blob", "err", err) + http.NotFound(w, r) + return + } + defer blob.Close() + + contentBytes, err := io.ReadAll(blob) + if err != nil { + l.Error("failed to read blob", "err", err) + } + content = string(contentBytes) + } - params = s.makeFileFragmentParams(&string, file.Name, content, forceCode) + params = s.makeFileFragmentParams(&str, file.Name, content, forceCode) } s.Pages.StringFileFragment(w, params) } @@ -653,7 +758,7 @@ func (s *Strings) getRecordCid(ctx context.Context, uri syntax.ATURI) (syntax.CI return "", err } - xrpcc := xrpc.Client{Host: ident.PDSEndpoint()} + xrpcc := indigoxrpc.Client{Host: ident.PDSEndpoint()} out, err := agnostic.RepoGetRecord(ctx, &xrpcc, "", uri.Collection().String(), ident.DID.String(), uri.RecordKey().String()) if err != nil { return "", err @@ -669,3 +774,11 @@ func (s *Strings) getRecordCid(ctx context.Context, uri syntax.ATURI) (syntax.CI return cid, nil } + +func gz(s string) io.Reader { + var b bytes.Buffer + w := gzip.NewWriter(&b) + w.Write([]byte(s)) + w.Close() + return &b +} diff --git a/blobstore/blobstore.go b/blobstore/blobstore.go new file mode 100644 index 00000000..fe1691c2 --- /dev/null +++ b/blobstore/blobstore.go @@ -0,0 +1,13 @@ +package blobstore + +import ( + "context" + "io" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/ipfs/go-cid" +) + +type BlobStore interface { + GetBlob(ctx context.Context, did syntax.DID, cid cid.Cid) (io.ReadCloser, error) +} diff --git a/blobstore/pds.go b/blobstore/pds.go new file mode 100644 index 00000000..3d413f14 --- /dev/null +++ b/blobstore/pds.go @@ -0,0 +1,46 @@ +package blobstore + +import ( + "context" + "fmt" + "io" + "net/http" + "net/url" + + "github.com/bluesky-social/indigo/atproto/identity" + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/ipfs/go-cid" +) + +type Pds struct { + dir identity.Directory +} + +func NewPdsBlobStore(dir identity.Directory) *Pds { + return &Pds{dir} +} + +func (s *Pds) GetBlob(ctx context.Context, did syntax.DID, cid cid.Cid) (io.ReadCloser, error) { + id, err := s.dir.LookupDID(ctx, did) + if err != nil { + return nil, err + } + + url, _ := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", id.PDSEndpoint())) + q := url.Query() + q.Set("did", did.String()) + q.Set("cid", cid.String()) + url.RawQuery = q.Encode() + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, url.String(), nil) + resp, err := http.DefaultClient.Do(req) + if err != nil { + return nil, err + } + + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("unexpected status: %s", resp.Status) + } + + return resp.Body, nil +} -- 2.51.2