From 3ba0b87560ee0598ce9abb2eae816890e7f93de7 Mon Sep 17 00:00:00 2001 From: Will Andrews Date: Sun, 22 Feb 2026 22:27:08 +0000 Subject: [PATCH] Run on cron job. Store blobs locally as a cache. --- .env.example | 1 + .gitignore | 1 + Dockerfile | 19 +++++++ docker-compose.yaml | 9 ++++ go.mod | 1 + go.sum | 2 + main.go | 45 +++++++++++++++- pds.go | 125 +++++++++++++++++++++++++------------------- readme.md | 3 +- tangled_knot.go | 20 +++---- 10 files changed, 161 insertions(+), 65 deletions(-) create mode 100644 Dockerfile create mode 100644 docker-compose.yaml diff --git a/.env.example b/.env.example index 7d66919..0a17398 100644 --- a/.env.example +++ b/.env.example @@ -7,3 +7,4 @@ PDS_HOST="https://your-pds.com" TANGLED_KNOT_DATABASE_DIRECTORY="/path/to/database/directory" TANGLED_KNOT_REPOSITORY_DIRECTORY="/path/to/repository/directory" BUGSNAG_API_KEY="enter-api-key-to-enable" +BLOB_DIR="./blobs" diff --git a/.gitignore b/.gitignore index 4c49bd7..3323b34 100644 --- a/.gitignore +++ b/.gitignore @@ -1 +1,2 @@ .env +.DS_Store diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..1f67951 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,19 @@ +FROM golang:alpine AS builder + +WORKDIR /app + +COPY . . +RUN go mod download + +COPY . . + +RUN CGO_ENABLED=0 go build -o back-at-it . + +FROM alpine:latest + +RUN apk --no-cache add ca-certificates + +WORKDIR /app/ +COPY --from=builder /app/back-at-it . + +ENTRYPOINT ["./back-at-it"] diff --git a/docker-compose.yaml b/docker-compose.yaml new file mode 100644 index 0000000..1176269 --- /dev/null +++ b/docker-compose.yaml @@ -0,0 +1,9 @@ +services: + back-at-it: + container_name: back-at-it + image: willdot/back-at-it:latest + environment: + ENV_LOCATION: "/app/data/ back-at-it.env" + volumes: + - ./data:/app/data + restart: always diff --git a/go.mod b/go.mod index 380cad4..b1b315d 100644 --- a/go.mod +++ b/go.mod @@ -21,6 +21,7 @@ require ( github.com/minio/md5-simd v1.1.2 // indirect github.com/philhofer/fwd v1.2.0 // indirect github.com/pkg/errors v0.9.1 // indirect + github.com/robfig/cron v1.2.0 // indirect github.com/rs/xid v1.6.0 // indirect github.com/stretchr/testify v1.10.0 // indirect github.com/tinylib/msgp v1.3.0 // indirect diff --git a/go.sum b/go.sum index da89b26..e1efb3a 100644 --- a/go.sum +++ b/go.sum @@ -34,6 +34,8 @@ github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/robfig/cron v1.2.0 h1:ZjScXvvxeQ63Dbyxy76Fj3AT3Ut0aKsyd2/tl3DTMuQ= +github.com/robfig/cron v1.2.0/go.mod h1:JGuDeoQd7Z6yL4zQhZ3OPEVHB7fL6Ka6skscFHfmt2k= github.com/rs/xid v1.6.0 h1:fV591PaemRlL6JfRxGDEPl69wICngIQ3shQtzfy2gxU= github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0= github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= diff --git a/main.go b/main.go index 36a4cc7..a0e619f 100644 --- a/main.go +++ b/main.go @@ -3,14 +3,26 @@ package main import ( "context" "log/slog" + "net/http" "os" + "time" "github.com/bugsnag/bugsnag-go/v2" "github.com/joho/godotenv" "github.com/minio/minio-go/v7" "github.com/minio/minio-go/v7/pkg/credentials" + "github.com/robfig/cron" ) +type service struct { + pdsHost string + did string + blobDir string + bucketName string + httpClient *http.Client + minioClient *minio.Client +} + func main() { ctx := context.Background() @@ -40,8 +52,37 @@ func main() { return } - backupPDS(ctx, minioClient, bucketName) - backupTangledKnot(ctx, minioClient, bucketName) + pdsHost := os.Getenv("PDS_HOST") + did := os.Getenv("DID") + blobDir := os.Getenv("BLOB_DIR") + + service := service{ + pdsHost: pdsHost, + did: did, + blobDir: blobDir, + bucketName: bucketName, + minioClient: minioClient, + httpClient: &http.Client{ + Timeout: time.Second * 5, + Transport: &http.Transport{ + IdleConnTimeout: time.Second * 90, + }, + }, + } + + service.backupPDS(ctx) + service.backupTangledKnot(ctx) + + c := cron.New() + + c.AddFunc("@hourly", func() { + service.backupPDS(ctx) + }) + c.AddFunc("@hourly", func() { + service.backupTangledKnot(ctx) + }) + + c.Start() } func createMinioClient() (*minio.Client, error) { diff --git a/pds.go b/pds.go index ed2e25d..cfd13e6 100644 --- a/pds.go +++ b/pds.go @@ -9,51 +9,52 @@ import ( "log/slog" "net/http" "os" + "path/filepath" + "time" "github.com/bugsnag/bugsnag-go/v2" "github.com/minio/minio-go/v7" ) -func backupPDS(ctx context.Context, minioClient *minio.Client, bucketName string) { - if os.Getenv("PDS_HOST") == "" || os.Getenv("DID") == "" { +func (s *service) backupPDS(ctx context.Context) { + if s.pdsHost == "" || s.did == "" { slog.Info("PDS_HOST or DID env not set - skipping PDS backup") return } - err := backupRepo(ctx, minioClient, bucketName) + err := s.backupRepo(ctx) if err != nil { slog.Error("backup repo", "error", err) bugsnag.Notify(err) return } - err = backupBlobs(ctx, minioClient, bucketName) + err = s.backupBlobs(ctx) if err != nil { slog.Error("backup blobs", "error", err) bugsnag.Notify(err) return } -} -func backupRepo(ctx context.Context, minioClient *minio.Client, bucketName string) error { - pdsHost := os.Getenv("PDS_HOST") - did := os.Getenv("DID") + slog.Info("finished PDS backup") +} - url := fmt.Sprintf("%s/xrpc/com.atproto.sync.getRepo?did=%s", pdsHost, did) +func (s *service) backupRepo(ctx context.Context) error { + url := fmt.Sprintf("%s/xrpc/com.atproto.sync.getRepo?did=%s", s.pdsHost, s.did) req, err := http.NewRequestWithContext(ctx, "GET", url, nil) if err != nil { return fmt.Errorf("create get repo request: %w", err) } req.Header.Add("ACCEPT", "application/vnd.ipld.car") - resp, err := http.DefaultClient.Do(req) + resp, err := s.httpClient.Do(req) if err != nil { return fmt.Errorf("get repo: %w", err) } defer resp.Body.Close() - _, err = minioClient.PutObject(ctx, bucketName, "pds-repo", resp.Body, -1, minio.PutObjectOptions{}) + _, err = s.minioClient.PutObject(ctx, s.bucketName, "pds-repo", resp.Body, -1, minio.PutObjectOptions{}) if err != nil { return fmt.Errorf("stream repo to bucket: %w", err) } @@ -61,44 +62,51 @@ func backupRepo(ctx context.Context, minioClient *minio.Client, bucketName strin return nil } -func backupBlobs(ctx context.Context, minioClient *minio.Client, bucketName string) error { - cids, err := getAllBlobCIDs(ctx) +func (s *service) backupBlobs(ctx context.Context) error { + cids, err := s.getAllBlobCIDs(ctx) if err != nil { return fmt.Errorf("get all blob CIDs: %w", err) } - reader, writer := io.Pipe() - defer reader.Close() + filename := fmt.Sprintf("%s-%s-blobs.zip", s.did, time.Now()) + defer os.Remove(filename) - zipWriter := zip.NewWriter(writer) + f, err := os.Create(filename) + if err != nil { + return fmt.Errorf("creating zip file: %w", err) + } + defer f.Close() - go func() { - defer writer.Close() - defer zipWriter.Close() + zipWriter := zip.NewWriter(f) + for _, cid := range cids { + slog.Info("processing cid", "cid", cid) + blob, err := s.getBlob(ctx, cid) + if err != nil { + slog.Error("failed to get blob", "cid", cid, "error", err) + bugsnag.Notify(err) + continue + } - for _, cid := range cids { - slog.Info("processing cid", "cid", cid) - blob, err := getBlob(ctx, cid) - if err != nil { - slog.Error("failed to get blob", "cid", cid, "error", err) - bugsnag.Notify(err) - continue - } + zipFile, err := zipWriter.Create(cid) + if err != nil { + slog.Error("create new file in zipwriter", "cid", cid, "error", err) + bugsnag.Notify(err) + continue + } - zipFile, err := zipWriter.Create(cid) - if err != nil { - slog.Error("create new file in zipwriter", "cid", cid, "error", err) - bugsnag.Notify(err) - blob.Close() - continue - } + zipFile.Write(blob) + } + err = zipWriter.Close() + if err != nil { + return fmt.Errorf("close zip writer: %w", err) + } - io.Copy(zipFile, blob) - blob.Close() - } - }() + fi, err := f.Stat() + if err != nil { + return fmt.Errorf("stat zip file: %w", err) + } - _, err = minioClient.PutObject(ctx, bucketName, "pds-blobs.zip", reader, -1, minio.PutObjectOptions{}) + _, err = s.minioClient.PutObject(ctx, s.bucketName, "pds-blobs.zip", f, fi.Size(), minio.PutObjectOptions{}) if err != nil { return fmt.Errorf("stream blobs to bucket: %w", err) } @@ -106,12 +114,12 @@ func backupBlobs(ctx context.Context, minioClient *minio.Client, bucketName stri return nil } -func getAllBlobCIDs(ctx context.Context) ([]string, error) { +func (s *service) getAllBlobCIDs(ctx context.Context) ([]string, error) { cursor := "" limit := 100 var cids []string for { - res, err := listBlobs(ctx, cursor, int64(limit)) + res, err := s.listBlobs(ctx, cursor, int64(limit)) if err != nil { return nil, fmt.Errorf("list blobs: %w", err) } @@ -134,18 +142,15 @@ type listBlobsResponse struct { CIDs []string `json:"cids"` } -func listBlobs(ctx context.Context, cursor string, limit int64) (listBlobsResponse, error) { - pdsHost := os.Getenv("PDS_HOST") - did := os.Getenv("DID") - +func (s *service) listBlobs(ctx context.Context, cursor string, limit int64) (listBlobsResponse, error) { // TODO: do proper url encoding of query params - url := fmt.Sprintf("%s/xrpc/com.atproto.sync.listBlobs?did=%s&cursor=%s&limit=%d", pdsHost, did, cursor, limit) + url := fmt.Sprintf("%s/xrpc/com.atproto.sync.listBlobs?did=%s&cursor=%s&limit=%d", s.pdsHost, s.did, cursor, limit) req, err := http.NewRequestWithContext(ctx, "GET", url, nil) if err != nil { return listBlobsResponse{}, fmt.Errorf("create list blobs request: %w", err) } - resp, err := http.DefaultClient.Do(req) + resp, err := s.httpClient.Do(req) if err != nil { return listBlobsResponse{}, fmt.Errorf("list blobs: %w", err) } @@ -166,21 +171,35 @@ func listBlobs(ctx context.Context, cursor string, limit int64) (listBlobsRespon return result, nil } -func getBlob(ctx context.Context, cid string) (io.ReadCloser, error) { - pdsHost := os.Getenv("PDS_HOST") - did := os.Getenv("DID") +func (s *service) getBlob(ctx context.Context, cid string) ([]byte, error) { + filename := filepath.Join(s.blobDir, fmt.Sprintf("blob-%s-%s", s.did, cid)) + existing, err := os.ReadFile(filename) + if !os.IsNotExist(err) { + return existing, nil + } // TODO: do proper url encoding of query params - url := fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob?did=%s&cid=%s", pdsHost, did, cid) + url := fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob?did=%s&cid=%s", s.pdsHost, s.did, cid) req, err := http.NewRequestWithContext(ctx, "GET", url, nil) if err != nil { return nil, fmt.Errorf("create get blob request: %w", err) } - resp, err := http.DefaultClient.Do(req) + resp, err := s.httpClient.Do(req) if err != nil { return nil, fmt.Errorf("get blob: %w", err) } + defer resp.Body.Close() + + b, err := io.ReadAll(resp.Body) + if err != nil { + return nil, fmt.Errorf("read blob response body: %w", err) + } + + err = os.WriteFile(filename, b, os.ModePerm) + if err != nil { + slog.Error("writing blob", "error", err, "cid", cid) + } - return resp.Body, nil + return b, nil } diff --git a/readme.md b/readme.md index b807c0f..032ab11 100644 --- a/readme.md +++ b/readme.md @@ -24,6 +24,7 @@ Run `go run .` ### Todo -- [ ] - Turn this into a long running app using a cron library perhaps +- [x] - Turn this into a long running app using a cron library perhaps +- [ ] - Work out how to tar just the directory requested, not the full path (ie `/home/will/tangled/repo` should back up `/repo` only and not create the full path of empty directories) - [ ] - User query params properly when creating the URLs to fetch repo and blobs - [ ] - Allow configuring the backup of knot repo data per users DID maybe? diff --git a/tangled_knot.go b/tangled_knot.go index a6ebf9a..7750355 100644 --- a/tangled_knot.go +++ b/tangled_knot.go @@ -13,12 +13,14 @@ import ( "github.com/minio/minio-go/v7" ) -func backupTangledKnot(ctx context.Context, minioClient *minio.Client, bucketName string) { - backupKnotDB(ctx, minioClient, bucketName) - backupKnotRepos(ctx, minioClient, bucketName) +func (s *service) backupTangledKnot(ctx context.Context) { + s.backupKnotDB(ctx) + s.backupKnotRepos(ctx) + + slog.Info("finished tangled knot backup") } -func backupKnotDB(ctx context.Context, minioClient *minio.Client, bucketName string) { +func (s *service) backupKnotDB(ctx context.Context) { dir := os.Getenv("TANGLED_KNOT_DATABASE_DIRECTORY") if dir == "" { slog.Info("TANGLED_KNOT_DATABASE_DIRECTORY env not set - skipping knot DB backup") @@ -29,14 +31,14 @@ func backupKnotDB(ctx context.Context, minioClient *minio.Client, bucketName str go compress(dir, pipeWriter) - _, err := minioClient.PutObject(ctx, bucketName, "knot-db.zip", pipeReader, -1, minio.PutObjectOptions{}) + _, err := s.minioClient.PutObject(ctx, s.bucketName, "knot-db.zip", pipeReader, -1, minio.PutObjectOptions{}) if err != nil { - slog.Error("stream knot DB to bucket: %w") + slog.Error("stream knot DB to bucket", "error", err) bugsnag.Notify(err) } } -func backupKnotRepos(ctx context.Context, minioClient *minio.Client, bucketName string) { +func (s *service) backupKnotRepos(ctx context.Context) { dir := os.Getenv("TANGLED_KNOT_REPOSITORY_DIRECTORY") if dir == "" { slog.Info("TANGLED_KNOT_REPOSITORY_DIRECTORY env not set - skipping knot repo backup") @@ -47,9 +49,9 @@ func backupKnotRepos(ctx context.Context, minioClient *minio.Client, bucketName go compress(dir, pipeWriter) - _, err := minioClient.PutObject(ctx, bucketName, "knot-repos.zip", pipeReader, -1, minio.PutObjectOptions{}) + _, err := s.minioClient.PutObject(ctx, s.bucketName, "knot-repos.zip", pipeReader, -1, minio.PutObjectOptions{}) if err != nil { - slog.Error("stream knot repos to bucket: %w") + slog.Error("stream knot repos to bucket", "error", err) bugsnag.Notify(err) } } -- 2.51.2