diff options
| author | Paul Buetow <paul@buetow.org> | 2024-09-21 14:03:45 +0300 |
|---|---|---|
| committer | Paul Buetow <paul@buetow.org> | 2024-09-21 14:03:45 +0300 |
| commit | dff4d455e07d639b82a0bed814f41d0656e9b6d0 (patch) | |
| tree | 56289a6fd80a00a9724992ab9cb8b3d7215ebedf /internal/server | |
| parent | 780ade3dc066afb8a43be824373414f3d316ebd6 (diff) | |
cleanup
Diffstat (limited to 'internal/server')
| -rw-r--r-- | internal/server/cron/cron.go | 45 | ||||
| -rw-r--r-- | internal/server/handler/handler.go | 72 | ||||
| -rw-r--r-- | internal/server/health/health.go | 107 | ||||
| -rw-r--r-- | internal/server/health/health_test.go | 24 | ||||
| -rw-r--r-- | internal/server/repository/pending.go | 57 | ||||
| -rw-r--r-- | internal/server/repository/pending_test.go | 66 | ||||
| -rw-r--r-- | internal/server/repository/repository.go | 320 | ||||
| -rw-r--r-- | internal/server/repository/repository_test.go | 363 | ||||
| -rw-r--r-- | internal/server/repository/stats.go | 22 | ||||
| -rw-r--r-- | internal/server/scheduler/scheduler.go | 16 | ||||
| -rw-r--r-- | internal/server/server.go | 49 |
11 files changed, 0 insertions, 1141 deletions
diff --git a/internal/server/cron/cron.go b/internal/server/cron/cron.go deleted file mode 100644 index 3d6b9cf..0000000 --- a/internal/server/cron/cron.go +++ /dev/null @@ -1,45 +0,0 @@ -package cron - -import ( - "context" - "log" - "time" - - config "codeberg.org/snonux/gos/internal/config/server" - "codeberg.org/snonux/gos/internal/server/health" - "codeberg.org/snonux/gos/internal/server/repository" - "codeberg.org/snonux/gos/internal/server/scheduler" -) - -func Run(ctx context.Context, conf config.ServerConfig, status health.Status) { - helloTicker := time.NewTicker(time.Hour) - mergeTicker := time.NewTicker(time.Second * time.Duration(conf.MergeIntervalS)) - scheduleTicker := time.NewTicker(time.Second * time.Duration(conf.ScheduleIntervalS)) - - for { - select { - case <-ctx.Done(): - return - case <-helloTicker.C: - run(ctx, "cron->Hello", status, func(ctx context.Context) error { - log.Println("hello world") - return nil - }) - case <-mergeTicker.C: - run(ctx, "cron->repository.Merge", status, repository.Instance(conf).MergeRemotely) - case <-scheduleTicker.C: - run(ctx, "cron->scheduler.Run", status, func(ctx context.Context) error { - return scheduler.Run(ctx, conf) - }) - } - } -} - -func run(ctx context.Context, what string, status health.Status, cb func(ctx context.Context) error) { - log.Println("CRON ticker initiating", what) - if err := cb(ctx); err != nil { - status.Set(health.Critical, what, err) - return - } - status.Clear(what) -} diff --git a/internal/server/handler/handler.go b/internal/server/handler/handler.go deleted file mode 100644 index a108f93..0000000 --- a/internal/server/handler/handler.go +++ /dev/null @@ -1,72 +0,0 @@ -package handler - -import ( - "context" - "fmt" - "io" - "net/http" - - "codeberg.org/snonux/gos/internal/config/server" - "codeberg.org/snonux/gos/internal/server/repository" - "codeberg.org/snonux/gos/internal/types" -) - -type Handler struct { - conf server.ServerConfig -} - -func New(conf server.ServerConfig) Handler { - return Handler{ - conf: conf, - } -} - -func (h Handler) Submit(ctx context.Context, w http.ResponseWriter, r *http.Request) error { - if r.Method != "POST" { - return fmt.Errorf("expected POST request, but got %s", r.Method) - } - - bytes, err := io.ReadAll(r.Body) - if err != nil { - return err - } - - entry, err := types.NewEntry(bytes) - if err != nil { - return err - } - return repository.Instance(h.conf).Merge(entry) -} - -func (h Handler) List(w http.ResponseWriter, r *http.Request) error { - if r.Method != "GET" { - return fmt.Errorf("expexted GET request") - } - - list, err := repository.Instance(h.conf).ListBytes() - if err != nil { - return err - } - - _, err = w.Write(list) - return err -} - -func (h Handler) Get(w http.ResponseWriter, r *http.Request) error { - json, err := repository.Instance(h.conf).GetJSON(r.URL.Query().Get("id")) - if err != nil { - return err - } - - fmt.Fprint(w, json) - return nil -} - -func (h Handler) Merge(ctx context.Context, w http.ResponseWriter, r *http.Request) error { - if err := repository.Instance(h.conf).MergeRemotely(ctx); err != nil { - return err - } - - fmt.Fprint(w, "Repository merge went well") - return nil -} diff --git a/internal/server/health/health.go b/internal/server/health/health.go deleted file mode 100644 index 5144416..0000000 --- a/internal/server/health/health.go +++ /dev/null @@ -1,107 +0,0 @@ -package health - -import ( - "fmt" - "log" - "strings" - "sync" -) - -type Severity int - -const ( - OK Severity = iota - Warning - Critical - Unknown -) - -func (s Severity) String() string { - switch s { - case OK: - return "OK" - case Warning: - return "WARNING" - case Critical: - return "CRITICAL" - case Unknown: - fallthrough - default: - return "UNKNOWN" - } -} - -type alert struct { - text string - severity Severity -} - -func (a alert) String() string { - return fmt.Sprintf("%s: %s", a.severity, a.text) -} - -type Status struct { - alerts map[string]alert - mu *sync.Mutex -} - -func NewStatus() Status { - return Status{ - alerts: make(map[string]alert), - mu: &sync.Mutex{}, - } -} - -func (hs Status) Set(s Severity, healthStatusKey string, info any) { - hs.mu.Lock() - defer hs.mu.Unlock() - - infoStr := fmt.Sprintf("%v", info) - log.Printf("status: alerting %s as %s: %s", healthStatusKey, s, infoStr) - - hs.alerts[healthStatusKey] = alert{ - text: infoStr, - severity: s, - } -} - -func (hs Status) Clear(healthStatusKey string) { - hs.mu.Lock() - defer hs.mu.Unlock() - - if _, ok := hs.alerts[healthStatusKey]; ok { - log.Println("status: clearing ", healthStatusKey) - delete(hs.alerts, healthStatusKey) - } -} - -func (hs Status) String() string { - var ( - alerts [4][]string // Alerts by severity - sb strings.Builder - ) - - hs.mu.Lock() - defer hs.mu.Unlock() - - for healthStatusKey, alert := range hs.alerts { - str := fmt.Sprintf("%s (handler %s)", alert, healthStatusKey) - alerts[alert.severity] = append(alerts[alert.severity], str) - } - - possible := [4]Severity{Unknown, Critical, Warning, OK} - for _, severity := range possible { - if len(alerts[severity]) == 0 { - continue - } - for _, alert := range alerts[severity] { - sb.WriteString(alert) - sb.WriteString("\n") - } - } - - if result := sb.String(); result != "" { - return result - } - return "OK: all is fine :-)\n" -} diff --git a/internal/server/health/health_test.go b/internal/server/health/health_test.go deleted file mode 100644 index 8d16d4b..0000000 --- a/internal/server/health/health_test.go +++ /dev/null @@ -1,24 +0,0 @@ -package health - -import "testing" - -func TestHealthStatus(t *testing.T) { - t.Parallel() - - hs := NewStatus() - hs.Set(Warning, "foo", "this is not good") - hs.Set(Critical, "bar", "this is not good either") - hs.Set(Warning, "baz", "urgh!") - hs.Set(Unknown, "baz", "don't know what happened here!") - hs.Clear("foo") - - result := hs.String() - expected := `UNKNOWN: don't know what happened here! (handler baz) -CRITICAL: this is not good either (handler bar) -` - - if result != expected { - t.Error("expected", expected, "but got", result) - } - t.Log("got as expexted", result) -} diff --git a/internal/server/repository/pending.go b/internal/server/repository/pending.go deleted file mode 100644 index a1fd7a7..0000000 --- a/internal/server/repository/pending.go +++ /dev/null @@ -1,57 +0,0 @@ -package repository - -import "codeberg.org/snonux/gos/internal/types" - -type pendingEntries map[types.EntryID]struct{} - -// Keep track of pending entries per social platform -type pending struct { - platforms map[types.PlatformName]pendingEntries -} - -func newPending() pending { - return pending{make(map[types.PlatformName]pendingEntries)} -} - -// Returns number of pending entries for the platform -// func (p pending) num(platform types.PlatformName) int { -// pe, ok := p.platforms[platform] -// if !ok { -// return 0 -// } -// return len(pe) -// } - -func (p pending) add(platform types.PlatformName, id types.EntryID) { - pe, ok := p.platforms[platform] - if !ok { - pe = make(pendingEntries) - } - pe[id] = struct{}{} - p.platforms[platform] = pe -} - -func (p pending) delete(platform types.PlatformName, id types.EntryID) { - pe, ok := p.platforms[platform] - if !ok { - return - } - delete(pe, id) - p.platforms[platform] = pe -} - -func (p pending) get(platform types.PlatformName) (pendingEntries, bool) { - pe, ok := p.platforms[platform] - return pe, ok && len(pe) > 0 -} - -func (p pending) next(platform types.PlatformName) (types.EntryID, bool) { - pe, ok := p.get(platform) - if !ok { - return "", false - } - for id := range pe { - return id, true - } - return "", false -} diff --git a/internal/server/repository/pending_test.go b/internal/server/repository/pending_test.go deleted file mode 100644 index 28563a6..0000000 --- a/internal/server/repository/pending_test.go +++ /dev/null @@ -1,66 +0,0 @@ -package repository - -import ( - "testing" - - "codeberg.org/snonux/gos/internal/types" -) - -func TestPendingAdd(t *testing.T) { - pending := newPending() - - entries, ok := pending.get(types.LinkedIn) - if ok { - t.Error("expected no ok return status") - } - if len(entries) != 0 { - t.Error("expected no entries") - } - - pending.add(types.LinkedIn, "fooid") - pending.add(types.LinkedIn, "barid") - - entries, ok = pending.get(types.LinkedIn) - if !ok { - t.Error("expected ok return status") - } - if len(entries) != 2 { - t.Error("expected two entries") - } -} - -func TestPendingDelete(t *testing.T) { - pending := newPending() - pending.add(types.LinkedIn, "fooid") - - entries, ok := pending.get(types.LinkedIn) - if !ok { - t.Error("expected ok return status") - } - if len(entries) != 1 { - t.Error("expected one entry") - } - - pending.delete(types.LinkedIn, "fooid") - if entries, ok = pending.get(types.LinkedIn); ok { - t.Error("expected not an ok", entries) - } -} - -func TestPendingNext(t *testing.T) { - pending := newPending() - - id, ok := pending.next(types.LinkedIn) - if ok { - t.Error("not expected ok return status", id) - } - - pending.add(types.LinkedIn, "fooid") - id, ok = pending.next(types.LinkedIn) - if !ok { - t.Error("expected ok return status") - } - if id != "fooid" { - t.Error("expected entry ID fooid") - } -} diff --git a/internal/server/repository/repository.go b/internal/server/repository/repository.go deleted file mode 100644 index d277884..0000000 --- a/internal/server/repository/repository.go +++ /dev/null @@ -1,320 +0,0 @@ -package repository - -import ( - "context" - "encoding/json" - "errors" - "fmt" - "log" - "regexp" - "sync" - "time" - - "codeberg.org/snonux/gos/internal/config/server" - "codeberg.org/snonux/gos/internal/easyhttp" - "codeberg.org/snonux/gos/internal/types" - "codeberg.org/snonux/gos/internal/vfs" -) - -var ( - instance Repository - once sync.Once -) - -type fs interface { - ReadFile(name string) ([]byte, error) - WriteFile(filePath string, bytes []byte) error - FindFiles(dataPath, suffix string) ([]string, error) -} - -// Contains an Entry ID and its checksumm, for the list and merge operations. -type entryPair struct { - ID, Checksum string -} - -// Holds all entries in the database / stores them to the disks.. -// TODO: Keep track of how many posts were made this week already. -type Repository struct { - pending pending - stats stats - conf server.ServerConfig - entries map[types.EntryID]types.Entry - mu *sync.Mutex - fs fs - loaded *bool - getIdRe *regexp.Regexp -} - -func Instance(conf server.ServerConfig) Repository { - once.Do(func() { - instance = newRepository(conf, vfs.RealFS{}) - }) - return instance -} - -// Need to register all social platforms for in-memory representation of shared posts and so on. -func newRepository(conf server.ServerConfig, fs fs) Repository { - var loaded bool - return Repository{ - pending: newPending(), // TODO: Make use of the pending for the selection algoritmh for the next post - stats: newStats(), // TODO: Make use of this. - conf: conf, - entries: make(map[types.EntryID]types.Entry), - mu: &sync.Mutex{}, - fs: fs, - loaded: &loaded, - getIdRe: regexp.MustCompile(`^[a-z0-9]{64}$`), - } -} - -// Gets next entry to be shared for the given social platform. -func (r Repository) Next(platform types.PlatformName) (types.Entry, bool) { - r.mu.Lock() - defer r.mu.Unlock() - - id, ok := r.pending.next(platform) - if !ok { - return types.Entry{}, false // No entry found - } - - var entry types.Entry - entry, ok = r.entries[id] - if !ok { - panic("did not expect that!") - } - return entry, true -} - -// Load repository into memory if not done yet. -func (r Repository) load() error { - if *r.loaded { - return nil - } - - filePaths, err := r.fs.FindFiles(r.conf.DataDir, ".json") - if err != nil { - return err - } - - var errs []error - for _, filePath := range filePaths { - log.Println("loading entry", filePath) - - bytes, err := r.fs.ReadFile(filePath) - if err != nil { - errs = append(errs, err) - continue - } - - entry, err := types.NewEntry(bytes) - if err != nil { - errs = append(errs, err) - continue - } - r.mu.Lock() - r.add(entry) - r.mu.Unlock() - } - - if len(errs) == 0 { - *r.loaded = true - } - - return errors.Join(errs...) -} - -func (r Repository) List() ([]entryPair, error) { - if err := r.load(); err != nil { - return []entryPair{}, err - } - - var pairs []entryPair - r.mu.Lock() - defer r.mu.Unlock() - - for _, entry := range r.entries { - pairs = append(pairs, entryPair{entry.ID, entry.Checksum()}) - } - - return pairs, nil -} - -func (r Repository) ListBytes() ([]byte, error) { - pairs, err := r.List() - if err != nil { - return []byte{}, err - } - return json.Marshal(pairs) -} - -func (r Repository) add(entry types.Entry) { - r.entries[entry.ID] = entry - - for _, platform := range r.conf.SocialPlatformsEnabled { - if entry.IsShared(platform) { - r.pending.delete(platform, entry.ID) - } else { - r.pending.add(platform, entry.ID) - } - } -} - -func (r Repository) persist(entry types.Entry) error { - r.add(entry) - - bytes, err := entry.JSONMarshal() - if err != err { - return err - } - return r.fs.WriteFile(r.entryPath(entry), bytes) -} - -func (r Repository) Get(id types.EntryID) (types.Entry, error) { - if !r.getIdRe.MatchString(id) { - return types.Entry{}, fmt.Errorf("invalid id %s", id) - } - if err := r.load(); err != nil { - return types.Entry{}, err - } - - r.mu.Lock() - defer r.mu.Unlock() - - entry, ok := r.entries[id] - if !ok { - return entry, fmt.Errorf("no entry with id %s found", id) - } - return entry, nil -} - -func (r Repository) GetJSON(id types.EntryID) (string, error) { - entry, err := r.Get(id) - if err != nil { - return "", err - } - - bytes, err := entry.JSONMarshal() - if err != nil { - return "", err - } - - return string(bytes), err -} - -func (r Repository) hasSameEntry(pair entryPair) bool { - r.mu.Lock() - defer r.mu.Unlock() - - entry, ok := r.entries[pair.ID] - if !ok || entry.Checksum() != pair.Checksum { - return false - } - return true -} - -func (r Repository) entryPath(ent types.Entry) string { - return fmt.Sprintf("%s/%s/%s.json", r.conf.DataDir, time.Now().Format("2006"), ent.ID) -} - -func (r Repository) Merge(otherEnt types.Entry) error { - if err := r.load(); err != nil { - return err - } - - r.mu.Lock() - defer r.mu.Unlock() - - entry, ok := r.entries[otherEnt.ID] - if !ok { - log.Println("can't find entry with ID", otherEnt.ID, "in local db, create new from copy") - var err error - if entry, err = types.NewEntryFromCopy(otherEnt); err != nil { - return err - } - return r.persist(entry) - } - - if entry, changed, err := entry.Update(otherEnt); changed { - if err != nil { - return err - } - return r.persist(entry) - } - return nil -} - -func (r Repository) MergeRemotely(ctx context.Context) error { - var errs []error - - if len(r.conf.Partners) == 0 { - log.Println("No partners configured - skipping remote merge operation") - return nil - } - - for _, partner := range r.conf.Partners { - if err := r.mergeRemotelyFromPartner(ctx, partner); err != nil { - errs = append(errs, err) - } - } - - return errors.Join(errs...) -} - -// Makes it mockable/testable -type getPairDataFunc func(context.Context, string, *[]entryPair) error -type getEntryDataFunc func(context.Context, string, string, *types.Entry) error - -func (r Repository) mergeRemotelyFromPartner(ctx context.Context, partner string) error { - getPair := func(ctx context.Context, partner string, pairs *[]entryPair) error { - uri := fmt.Sprintf("%s/list", partner) - return easyhttp.GetData(ctx, uri, r.conf.APIKey, pairs) - } - - getEntry := func(ctx context.Context, partner, id types.EntryID, entry *types.Entry) error { - uri := fmt.Sprintf("%s/get?id=%s", partner, id) - return easyhttp.GetData(ctx, uri, r.conf.APIKey, entry) - } - - return r.mergeFromPartner(ctx, partner, getPair, getEntry) -} - -func (r Repository) mergeFromPartner(ctx context.Context, partner string, - getPair getPairDataFunc, getEntry getEntryDataFunc) error { - - if err := r.load(); err != nil { - return err - } - - var ( - errs []error - pairs []entryPair - ) - - if err := getPair(ctx, partner, &pairs); err != nil { - return err - } - - for _, pair := range pairs { - if r.hasSameEntry(pair) { - continue - } - - log.Println("pair", pair, "missing in local reposotory, going to merge it") - - var entry types.Entry - if err := getEntry(ctx, partner, pair.ID, &entry); err != nil { - errs = append(errs, err) - continue - } - - // In theory, this should never happen - if pair.ID != entry.ID { - errs = append(errs, fmt.Errorf("pair ID %s does not match entry id %s", pair.ID, entry.ID)) - continue - } - - errs = append(errs, r.Merge(entry)) - } - - return errors.Join(errs...) -} diff --git a/internal/server/repository/repository_test.go b/internal/server/repository/repository_test.go deleted file mode 100644 index cdde29c..0000000 --- a/internal/server/repository/repository_test.go +++ /dev/null @@ -1,363 +0,0 @@ -package repository - -import ( - "context" - "fmt" - "testing" - - "codeberg.org/snonux/gos/internal/config/server" - "codeberg.org/snonux/gos/internal/types" - "codeberg.org/snonux/gos/internal/vfs" -) - -func TestRepositoryPutGet(t *testing.T) { - t.Parallel() - - fs := make(vfs.MemoryFS) - repo := newRepository(server.ServerConfig{DataDir: "./data"}, fs) - - for _, entry := range makeEntries(t) { - t.Run(entry.ID, func(t *testing.T) { - _ = repo.persist(entry) - entGot, err := repo.Get(entry.ID) - if err != nil { - t.Error(err) - } - if !entGot.Equals(entry) { - t.Error("expected to get", entry, "but got", entGot) - } - }) - } -} - -func TestRepositoryLoad(t *testing.T) { - t.Parallel() - - fs := make(vfs.MemoryFS) - repo := newRepository(server.ServerConfig{DataDir: "./data"}, fs) - entries := makeEntries(t) - - // Write entries into the VFS - for _, entry := range entries { - bytes, _ := entry.JSONMarshal() - _ = repo.fs.WriteFile(repo.entryPath(entry), bytes) - } - - // Load entries from VFS into the repo - if err := repo.load(); err != nil { - t.Error(err) - } - - for _, entry := range entries { - t.Run(entry.ID, func(t *testing.T) { - entGot, err := repo.Get(entry.ID) - if err != nil { - t.Error(err) - } - if !entGot.Equals(entry) { - t.Error("expected to get", entry, "but got", entGot) - } - }) - } -} - -func TestRepositoryList(t *testing.T) { - t.Parallel() - - fs := make(vfs.MemoryFS) - repo := newRepository(server.ServerConfig{DataDir: "./data"}, fs) - entries := makeEntries(t) - - for _, entry := range entries { - _ = repo.persist(entry) - } - - pairs, _ := repo.List() - if len(entries) != len(pairs) { - t.Error("expected as many entries as pairs") - } - - for _, entry := range entries { - var found bool - for _, pair := range pairs { - if entry.ID == pair.ID && entry.Checksum() == pair.Checksum { - found = true - t.Log("entry matches pair", entry, pair) - break - } - } - if !found { - t.Error("could not find entry", entry, "in", pairs) - } - } -} - -func TestRepositoryHasSameEntry(t *testing.T) { - t.Parallel() - - fs := make(vfs.MemoryFS) - repo := newRepository(server.ServerConfig{DataDir: "./data"}, fs) - entry, _ := makeAnEntry() - _ = repo.persist(entry) - - pair := entryPair{entry.ID, entry.Checksum()} - if !repo.hasSameEntry(pair) { - t.Error("repo does not contain entry corresponding to pair", pair) - } - - pair = entryPair{"nonexistent", "nonexistent"} - if repo.hasSameEntry(pair) { - t.Error("repo does contain entry corresponding to pair", pair, "but that should not be") - } -} - -func TestRepositoryMerge(t *testing.T) { - t.Parallel() - - fs := make(vfs.MemoryFS) - repo := newRepository(server.ServerConfig{DataDir: "./data"}, fs) - entry1, _ := makeAnEntry() - _ = repo.persist(entry1) - - entry2, _ := makeAnotherEntry() - // Need to have the same IDs so that the entries will actually be merged - entry2.ID = entry1.ID - // Merge a modified entry2 into the repository. - entry2.Body = "merged" - entry2.Epoch = 12345 - _ = repo.Merge(entry2) - - pairs, _ := repo.List() - // Ensuring the merge didn't add a new entry - if len(pairs) != 1 { - t.Error("expected exactly one element in the repo but got", pairs) - } - - entGot, _ := repo.Get(entry1.ID) - if entGot.Body != "merged" { - t.Error("unexpected body", entGot.Body) - } - if entGot.Epoch != 12345 { - t.Error("unexpected epoch", entGot.Epoch) - } -} - -func TestRepositoryMergeFromPartner(t *testing.T) { - fs1 := make(vfs.MemoryFS) - repo1 := newRepository(server.ServerConfig{DataDir: "./data1"}, fs1) - fs2 := make(vfs.MemoryFS) - repo2 := newRepository(server.ServerConfig{DataDir: "./data2"}, fs2) - - entry1, _ := makeAnEntry() - _ = repo1.persist(entry1) - entry2, _ := makeAnotherEntry() - _ = repo2.persist(entry2) - - getPair := func(ctx context.Context, partner string, pairs *[]entryPair) error { - var ( - pairs_ []entryPair - err error - ) - - switch partner { - case "repo1": - pairs_, err = repo1.List() - case "repo2": - pairs_, err = repo2.List() - } - - if err != nil { - return err - } - *pairs = pairs_ - - t.Log("got pairs", *pairs, "from repo", partner) - return nil - } - - getEntry := func(ctx context.Context, partner, id string, entry *types.Entry) error { - var ( - entry_ types.Entry - err error - ) - - switch partner { - case "repo1": - entry_, err = repo1.Get(id) - case "repo2": - entry_, err = repo2.Get(id) - } - - if err != nil { - return err - } - *entry = entry_ - - t.Log("got entry", *entry, "from repo", partner) - return nil - } - - // Compare both repos, they should now contain the same entries - compare := func(repo1, repo2 Repository) error { - pairs, err := repo1.List() - if err != nil { - return err - } - - for _, pair := range pairs { - entry1, err := repo1.Get(pair.ID) - if err != nil { - return err - } - entry2, err := repo2.Get(pair.ID) - if err != nil { - return err - } - - t.Log("comparing entries") - t.Log("entry1", entry1) - t.Log("entry2", entry2) - - if !entry1.Equals(entry2) { - return fmt.Errorf("entries entry1 and entry2 don't equal") - } - } - - return nil - } - - t.Run("Merge entries from repo2 into repo1", func(t *testing.T) { - if err := repo1.mergeFromPartner(context.Background(), "repo2", getPair, getEntry); err != nil { - t.Error(err) - } - if err := compare(repo2, repo1); err != nil { - t.Error(err) - } - }) - - t.Run("Merge entries from repo1 into repo2", func(t *testing.T) { - if err := repo2.mergeFromPartner(context.Background(), "repo1", getPair, getEntry); err != nil { - t.Error(err) - } - if err := compare(repo1, repo2); err != nil { - t.Error(err) - } - }) - - t.Run("Change shared flag and merge to partner", func(t *testing.T) { - entry, err := repo1.Get(entry1.ID) - if err != nil { - t.Error(err) - } - - // Validate the correct test setup - if entry.IsShared(types.LinkedIn) { - t.Error("for the test expected LinkedIn not to be shared") - } - - // Simulate that the entry was shared to LinkedIn social media! - linkedIn, ok := entry.Shared[types.LinkedIn] - if !ok { - t.Error("expected to have a LinkedIn shared entry") - } - linkedIn.Is = true - entry.Shared[types.LinkedIn] = linkedIn - - if err := repo1.Merge(entry); err != nil { - t.Error(err) - } - - // Before merging, repos should be out of sync. - if err := compare(repo1, repo2); err == nil { - t.Log("as expected repos are out of sync", err) - } - - // Partner is merging the repo. - if err := repo1.mergeFromPartner(context.Background(), "repo2", getPair, getEntry); err != nil { - t.Error(err) - } - - // Still out of sync, as we merged the repos the wrong direction. - if err := compare(repo1, repo2); err == nil { - t.Log("as expected repos are out of sync", err) - } - - // Partner is merging the repo the right direction. - if err := repo2.mergeFromPartner(context.Background(), "repo1", getPair, getEntry); err != nil { - t.Error(err) - } - - // Now, partners should be in sync. - if err := compare(repo1, repo2); err != nil { - t.Error(err) - } - }) -} - -func TestRepositoryNext(t *testing.T) { - t.Parallel() - - fs := make(vfs.MemoryFS) - repo := newRepository(server.ServerConfig{ - DataDir: "./data", - SocialPlatformsEnabled: []types.PlatformName{ - types.LinkedIn, types.Mastodon, types.Textfile, - }, - }, fs) - entries := makeEntries(t) - - for _, entry := range entries { - _ = repo.persist(entry) - } - - if entry, ok := repo.Next(types.Mastodon); ok { - t.Error("expected no Mastodon entry to be found", entry) - } - - if _, ok := repo.Next(types.LinkedIn); !ok { - t.Error("expected an unshared LinkedIn entry to be found") - } - - if _, ok := repo.Next(types.Textfile); !ok { - t.Error("expected an unshared Textfile entry to be found") - } -} - -func makeEntries(t *testing.T) []types.Entry { - entry1, err := makeAnEntry() - if err != nil { - t.Error(err) - } - entry2, err := makeAnotherEntry() - if err != nil { - t.Error(err) - } - return []types.Entry{entry1, entry2} -} - -func makeAnEntry() (types.Entry, error) { - entry := ` - { - "body": "Body text here", - "shared": { - "Mastodon": { "is": true }, - "LinkedIn": { "is": false } - } - } - ` - return types.NewEntry([]byte(entry)) -} - -func makeAnotherEntry() (types.Entry, error) { - entry := ` - { - "body": "Another text here", - "shared": { - "Mastodon": { "is": true }, - "LinkedIn": { "is": true }, - "Textfile": { "is": false } - } - } - ` - return types.NewEntry([]byte(entry)) -} diff --git a/internal/server/repository/stats.go b/internal/server/repository/stats.go deleted file mode 100644 index 005ef80..0000000 --- a/internal/server/repository/stats.go +++ /dev/null @@ -1,22 +0,0 @@ -package repository - -import "codeberg.org/snonux/gos/internal/types" - -// Keeps track of how many messages were posted to social media over the last week and month. -type stats struct { - // Sliding window of entries shared last 7 days - last7Days map[types.PlatformName][]types.UnixEpoch - // Sliding window of entries shared last 30 days - last30Days map[types.PlatformName][]types.UnixEpoch -} - -func newStats() stats { - return stats{ - last7Days: make(map[types.PlatformName][]types.UnixEpoch), - last30Days: make(map[types.PlatformName][]types.UnixEpoch), - } -} - -// func (s stats) add(platform types.PlatformName, entry types.Entry) { - -// } diff --git a/internal/server/scheduler/scheduler.go b/internal/server/scheduler/scheduler.go deleted file mode 100644 index bbcb59a..0000000 --- a/internal/server/scheduler/scheduler.go +++ /dev/null @@ -1,16 +0,0 @@ -package scheduler - -import ( - "context" - "log" - - "codeberg.org/snonux/gos/internal/config/server" -) - -// TODO: Finish implementing this -func Run(ctx context.Context, config server.ServerConfig) error { - for _, platform := range config.SocialPlatformsEnabled { - log.Println("TODO: implement ... posting a post now or what on", platform) - } - return nil -} diff --git a/internal/server/server.go b/internal/server/server.go deleted file mode 100644 index b2cb0d0..0000000 --- a/internal/server/server.go +++ /dev/null @@ -1,49 +0,0 @@ -package server - -import ( - "fmt" - "log" - "net/http" - - config "codeberg.org/snonux/gos/internal/config/server" - "codeberg.org/snonux/gos/internal/server/health" -) - -type Server struct { - Status health.Status - Conf config.ServerConfig -} - -type HandlerFuncWithError func(http.ResponseWriter, *http.Request) error - -func New(conf config.ServerConfig, status health.Status) Server { - return Server{Conf: conf, Status: status} -} - -func (serv Server) Handle(name string, handler HandlerFuncWithError) { - var ( - handlerPath = fmt.Sprintf("/%s", name) - handlerName = fmt.Sprintf("%sHandler", name) - ) - - http.HandleFunc(handlerPath, func(w http.ResponseWriter, r *http.Request) { - log.Println("Someone requested", handlerName) - - // The health endpoint doesn't require an API key - if handlerNa |
