feat: ticketing sync workers with idempotent import and orphan flagging (phase 6)
- task domain model + exhaustive §4.4 status machine tests (single source of truth) - per-customer polling workers reconciled every 30s from the customers collection; panic-safe runs; manual sync-now trigger; status provider for the admin panel - §5.3 upsert keyed on (system, key, customerId): refreshes content while status=imported, records upstream_changed timeline + notification after - attachments cached into GridFS via connector-authorized downloads - orphaned tickets flagged with timeline + consultant notification, never deleted; idempotent across polls - consultant identity from users.extra.ticketingIdentities per system - notifications store (insert/list/unread-count/mark-read) Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,79 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"go.mongodb.org/mongo-driver/v2/bson"
|
||||
"go.mongodb.org/mongo-driver/v2/mongo/options"
|
||||
|
||||
"bountyboard/internal/ulid"
|
||||
)
|
||||
|
||||
// Notification is an in-app notification (§4.8).
|
||||
type Notification struct {
|
||||
ID string `bson:"_id" json:"id"`
|
||||
UserID string `bson:"userId" json:"userId"`
|
||||
Kind string `bson:"kind" json:"kind"`
|
||||
Title string `bson:"title" json:"title"`
|
||||
Body string `bson:"body" json:"body"`
|
||||
Link string `bson:"link" json:"link"`
|
||||
ReadAt *time.Time `bson:"readAt" json:"readAt"`
|
||||
CreatedAt time.Time `bson:"createdAt" json:"createdAt"`
|
||||
}
|
||||
|
||||
func (s *Store) InsertNotification(ctx context.Context, n *Notification) error {
|
||||
n.ID = ulid.New()
|
||||
n.CreatedAt = time.Now().UTC()
|
||||
if _, err := s.DB.Collection("notifications").InsertOne(ctx, n); err != nil {
|
||||
return fmt.Errorf("insert notification: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Store) ListNotifications(ctx context.Context, userID string, unreadOnly bool, cursor string, limit int) ([]Notification, error) {
|
||||
if limit <= 0 || limit > 100 {
|
||||
limit = 30
|
||||
}
|
||||
q := bson.M{"userId": userID}
|
||||
if unreadOnly {
|
||||
q["readAt"] = nil
|
||||
}
|
||||
if cursor != "" {
|
||||
q["_id"] = bson.M{"$lt": cursor}
|
||||
}
|
||||
cur, err := s.DB.Collection("notifications").Find(ctx, q,
|
||||
options.Find().SetSort(bson.D{{Key: "_id", Value: -1}}).SetLimit(int64(limit)))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("list notifications: %w", err)
|
||||
}
|
||||
out := []Notification{}
|
||||
if err := cur.All(ctx, &out); err != nil {
|
||||
return nil, fmt.Errorf("decode notifications: %w", err)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (s *Store) CountUnreadNotifications(ctx context.Context, userID string) (int64, error) {
|
||||
n, err := s.DB.Collection("notifications").CountDocuments(ctx,
|
||||
bson.M{"userId": userID, "readAt": nil})
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("count unread: %w", err)
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
|
||||
// MarkNotificationsRead marks the given ids (or all when ids is empty) read.
|
||||
func (s *Store) MarkNotificationsRead(ctx context.Context, userID string, ids []string) (int64, error) {
|
||||
q := bson.M{"userId": userID, "readAt": nil}
|
||||
if len(ids) > 0 {
|
||||
q["_id"] = bson.M{"$in": ids}
|
||||
}
|
||||
res, err := s.DB.Collection("notifications").UpdateMany(ctx, q,
|
||||
bson.M{"$set": bson.M{"readAt": time.Now().UTC()}})
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("mark read: %w", err)
|
||||
}
|
||||
return res.ModifiedCount, nil
|
||||
}
|
||||
@@ -0,0 +1,330 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"go.mongodb.org/mongo-driver/v2/bson"
|
||||
"go.mongodb.org/mongo-driver/v2/mongo"
|
||||
"go.mongodb.org/mongo-driver/v2/mongo/options"
|
||||
|
||||
"bountyboard/internal/domain"
|
||||
"bountyboard/internal/ulid"
|
||||
)
|
||||
|
||||
func (s *Store) tasks() *mongo.Collection { return s.DB.Collection("tasks") }
|
||||
|
||||
func (s *Store) InsertTask(ctx context.Context, t *domain.Task) error {
|
||||
now := time.Now().UTC()
|
||||
if t.ID == "" {
|
||||
t.ID = ulid.New()
|
||||
}
|
||||
if t.RootID == "" {
|
||||
t.RootID = t.ID
|
||||
}
|
||||
t.CreatedAt, t.UpdatedAt, t.Version = now, now, 1
|
||||
if t.AcceptanceCriteria == nil {
|
||||
t.AcceptanceCriteria = []string{}
|
||||
}
|
||||
if t.Attachments == nil {
|
||||
t.Attachments = []domain.TaskAttachment{}
|
||||
}
|
||||
if t.Links == nil {
|
||||
t.Links = []string{}
|
||||
}
|
||||
if t.ClaimRequests == nil {
|
||||
t.ClaimRequests = []domain.ClaimRequest{}
|
||||
}
|
||||
if t.Timeline == nil {
|
||||
t.Timeline = []domain.TimelineEntry{}
|
||||
}
|
||||
if t.Comments == nil {
|
||||
t.Comments = []domain.Comment{}
|
||||
}
|
||||
if t.TimeLog == nil {
|
||||
t.TimeLog = []domain.TimeLogEntry{}
|
||||
}
|
||||
if _, err := s.tasks().InsertOne(ctx, t); err != nil {
|
||||
if mongo.IsDuplicateKeyError(err) {
|
||||
return fmt.Errorf("task %s: %w", t.ID, ErrDuplicate)
|
||||
}
|
||||
return fmt.Errorf("insert task: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Store) TaskByID(ctx context.Context, id string) (*domain.Task, error) {
|
||||
var t domain.Task
|
||||
err := s.tasks().FindOne(ctx, bson.M{"_id": id}).Decode(&t)
|
||||
if errors.Is(err, mongo.ErrNoDocuments) {
|
||||
return nil, ErrNotFound
|
||||
}
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("find task: %w", err)
|
||||
}
|
||||
return &t, nil
|
||||
}
|
||||
|
||||
// TaskFilter selects tasks for boards and workers.
|
||||
type TaskFilter struct {
|
||||
CustomerID string
|
||||
CustomerIDs []string
|
||||
ConsultantID string
|
||||
ConsultantIDs []string
|
||||
Statuses []domain.TaskStatus
|
||||
AssigneeID string
|
||||
RootID string
|
||||
ParentID string
|
||||
Search string // Mongo text index
|
||||
MinBounty float64
|
||||
Sort string // newest | bounty | updated
|
||||
Cursor string
|
||||
Limit int
|
||||
}
|
||||
|
||||
func (f TaskFilter) build() bson.M {
|
||||
q := bson.M{}
|
||||
if f.CustomerID != "" {
|
||||
q["customerId"] = f.CustomerID
|
||||
}
|
||||
if len(f.CustomerIDs) > 0 {
|
||||
q["customerId"] = bson.M{"$in": f.CustomerIDs}
|
||||
}
|
||||
if f.ConsultantID != "" {
|
||||
q["consultantId"] = f.ConsultantID
|
||||
}
|
||||
if len(f.ConsultantIDs) > 0 {
|
||||
q["consultantId"] = bson.M{"$in": f.ConsultantIDs}
|
||||
}
|
||||
if len(f.Statuses) > 0 {
|
||||
q["status"] = bson.M{"$in": f.Statuses}
|
||||
}
|
||||
if f.AssigneeID != "" {
|
||||
q["assignee.userId"] = f.AssigneeID
|
||||
}
|
||||
if f.RootID != "" {
|
||||
q["rootId"] = f.RootID
|
||||
}
|
||||
if f.ParentID != "" {
|
||||
q["parentId"] = f.ParentID
|
||||
}
|
||||
if f.Search != "" {
|
||||
q["$text"] = bson.M{"$search": f.Search}
|
||||
}
|
||||
if f.MinBounty > 0 {
|
||||
q["bounty"] = bson.M{"$gte": f.MinBounty}
|
||||
}
|
||||
return q
|
||||
}
|
||||
|
||||
func (s *Store) ListTasks(ctx context.Context, f TaskFilter) ([]domain.Task, error) {
|
||||
q := f.build()
|
||||
limit := f.Limit
|
||||
if limit <= 0 || limit > 500 {
|
||||
limit = 100
|
||||
}
|
||||
sort := bson.D{{Key: "_id", Value: -1}} // newest first (ULID)
|
||||
switch f.Sort {
|
||||
case "bounty":
|
||||
sort = bson.D{{Key: "bounty", Value: -1}, {Key: "_id", Value: -1}}
|
||||
case "updated":
|
||||
sort = bson.D{{Key: "updatedAt", Value: -1}}
|
||||
}
|
||||
if f.Cursor != "" && f.Sort != "bounty" && f.Sort != "updated" {
|
||||
q["_id"] = bson.M{"$lt": f.Cursor}
|
||||
}
|
||||
cur, err := s.tasks().Find(ctx, q, options.Find().SetSort(sort).SetLimit(int64(limit)))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("list tasks: %w", err)
|
||||
}
|
||||
out := []domain.Task{}
|
||||
if err := cur.All(ctx, &out); err != nil {
|
||||
return nil, fmt.Errorf("decode tasks: %w", err)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// UpsertResult describes what an imported-ticket upsert did.
|
||||
type UpsertResult struct {
|
||||
Created bool
|
||||
UpstreamChanged bool // changed after atomization began (timeline-only)
|
||||
Task *domain.Task
|
||||
}
|
||||
|
||||
// UpsertImportedTask implements the §5.3 idempotent upsert keyed on
|
||||
// (external.system, external.key, customerId). While the task is still
|
||||
// `imported`, upstream refreshes title/description/attachments/raw; after
|
||||
// atomization begins, a change only appends an "upstream ticket changed"
|
||||
// timeline entry (caller notifies the consultant).
|
||||
func (s *Store) UpsertImportedTask(ctx context.Context, incoming *domain.Task) (*UpsertResult, error) {
|
||||
if incoming.External == nil {
|
||||
return nil, errors.New("imported task requires external ref")
|
||||
}
|
||||
key := bson.M{
|
||||
"external.system": incoming.External.System,
|
||||
"external.key": incoming.External.Key,
|
||||
"customerId": incoming.CustomerID,
|
||||
}
|
||||
var existing domain.Task
|
||||
err := s.tasks().FindOne(ctx, key).Decode(&existing)
|
||||
switch {
|
||||
case errors.Is(err, mongo.ErrNoDocuments):
|
||||
incoming.Origin = domain.OriginImported
|
||||
incoming.Status = domain.StatusImported
|
||||
incoming.Timeline = []domain.TimelineEntry{{
|
||||
At: time.Now().UTC(), ActorID: "system", Event: "imported",
|
||||
Data: map[string]any{"key": incoming.External.Key},
|
||||
}}
|
||||
if err := s.InsertTask(ctx, incoming); err != nil {
|
||||
// A concurrent worker may have inserted the same key (unique
|
||||
// index): treat as found-and-unchanged.
|
||||
if errors.Is(err, ErrDuplicate) {
|
||||
return &UpsertResult{Created: false, Task: incoming}, nil
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
return &UpsertResult{Created: true, Task: incoming}, nil
|
||||
case err != nil:
|
||||
return nil, fmt.Errorf("find imported task: %w", err)
|
||||
}
|
||||
|
||||
if existing.External.ContentHash == incoming.External.ContentHash {
|
||||
return &UpsertResult{Task: &existing}, nil // nothing new upstream
|
||||
}
|
||||
|
||||
now := time.Now().UTC()
|
||||
if existing.Status == domain.StatusImported {
|
||||
_, err := s.tasks().UpdateOne(ctx, bson.M{"_id": existing.ID}, bson.M{
|
||||
"$set": bson.M{
|
||||
"title": incoming.Title,
|
||||
"description": incoming.Description,
|
||||
"attachments": incoming.Attachments,
|
||||
"external.raw": incoming.External.Raw,
|
||||
"external.url": incoming.External.URL,
|
||||
"external.type": incoming.External.Type,
|
||||
"external.contentHash": incoming.External.ContentHash,
|
||||
"updatedAt": now,
|
||||
},
|
||||
"$inc": bson.M{"version": 1},
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("refresh imported task: %w", err)
|
||||
}
|
||||
return &UpsertResult{Task: &existing}, nil
|
||||
}
|
||||
|
||||
// Atomization already began: record only.
|
||||
_, err = s.tasks().UpdateOne(ctx, bson.M{"_id": existing.ID}, bson.M{
|
||||
"$set": bson.M{"external.contentHash": incoming.External.ContentHash, "updatedAt": now},
|
||||
"$push": bson.M{"timeline": domain.TimelineEntry{
|
||||
At: now, ActorID: "system", Event: "upstream_changed",
|
||||
Data: map[string]any{"key": incoming.External.Key},
|
||||
}},
|
||||
"$inc": bson.M{"version": 1},
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("record upstream change: %w", err)
|
||||
}
|
||||
return &UpsertResult{UpstreamChanged: true, Task: &existing}, nil
|
||||
}
|
||||
|
||||
// MarkTaskOrphaned flags a task whose upstream ticket disappeared (§5.3).
|
||||
func (s *Store) MarkTaskOrphaned(ctx context.Context, id string) error {
|
||||
now := time.Now().UTC()
|
||||
res, err := s.tasks().UpdateOne(ctx,
|
||||
bson.M{"_id": id, "external.orphaned": bson.M{"$ne": true}},
|
||||
bson.M{
|
||||
"$set": bson.M{"external.orphaned": true, "updatedAt": now},
|
||||
"$push": bson.M{"timeline": domain.TimelineEntry{
|
||||
At: now, ActorID: "system", Event: "orphaned",
|
||||
Data: map[string]any{"reason": "no longer returned by upstream query"},
|
||||
}},
|
||||
"$inc": bson.M{"version": 1},
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("mark orphaned: %w", err)
|
||||
}
|
||||
if res.MatchedCount == 0 {
|
||||
return ErrNotFound // already orphaned or missing
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// AppendTimeline pushes a timeline entry without status change.
|
||||
func (s *Store) AppendTimeline(ctx context.Context, taskID string, e domain.TimelineEntry) error {
|
||||
if e.At.IsZero() {
|
||||
e.At = time.Now().UTC()
|
||||
}
|
||||
_, err := s.tasks().UpdateOne(ctx, bson.M{"_id": taskID}, bson.M{
|
||||
"$push": bson.M{"timeline": e},
|
||||
"$set": bson.M{"updatedAt": time.Now().UTC()},
|
||||
"$inc": bson.M{"version": 1},
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("append timeline: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// TransitionTask performs a version-checked §4.4 transition with a timeline
|
||||
// entry and optional extra $set fields.
|
||||
func (s *Store) TransitionTask(ctx context.Context, t *domain.Task, to domain.TaskStatus,
|
||||
actorID, event string, data map[string]any, extraSet bson.M) error {
|
||||
if err := domain.Transition(t.Status, to); err != nil {
|
||||
return err
|
||||
}
|
||||
set := bson.M{"status": to}
|
||||
for k, v := range extraSet {
|
||||
set[k] = v
|
||||
}
|
||||
update := bson.M{
|
||||
"$set": set,
|
||||
"$push": bson.M{"timeline": domain.TimelineEntry{
|
||||
At: time.Now().UTC(), ActorID: actorID, Event: event, Data: data,
|
||||
}},
|
||||
}
|
||||
return UpdateVersioned(ctx, s.tasks(), t.ID, t.Version, update)
|
||||
}
|
||||
|
||||
// SetTaskAttachments persists attachment metadata (e.g. cached fileIds).
|
||||
func (s *Store) SetTaskAttachments(ctx context.Context, taskID string, atts []domain.TaskAttachment) error {
|
||||
_, err := s.tasks().UpdateOne(ctx, bson.M{"_id": taskID}, bson.M{
|
||||
"$set": bson.M{"attachments": atts, "updatedAt": time.Now().UTC()},
|
||||
"$inc": bson.M{"version": 1},
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("set attachments: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ListImportedTasksForOrphanCheck returns imported-origin tasks for one
|
||||
// customer/system/consultant that are still live (not archived).
|
||||
func (s *Store) ListImportedTasksForOrphanCheck(ctx context.Context, customerID, system, consultantID string) ([]domain.Task, error) {
|
||||
cur, err := s.tasks().Find(ctx, bson.M{
|
||||
"customerId": customerID,
|
||||
"consultantId": consultantID,
|
||||
"origin": domain.OriginImported,
|
||||
"external.system": system,
|
||||
"status": bson.M{"$ne": domain.StatusArchived},
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("list for orphan check: %w", err)
|
||||
}
|
||||
out := []domain.Task{}
|
||||
if err := cur.All(ctx, &out); err != nil {
|
||||
return nil, fmt.Errorf("decode orphan check: %w", err)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// DeleteTasks removes tasks matching ids (re-run subdivision cleanup).
|
||||
func (s *Store) DeleteTasks(ctx context.Context, ids []string) (int64, error) {
|
||||
res, err := s.tasks().DeleteMany(ctx, bson.M{"_id": bson.M{"$in": ids}})
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("delete tasks: %w", err)
|
||||
}
|
||||
return res.DeletedCount, nil
|
||||
}
|
||||
@@ -115,6 +115,23 @@ func (s *Store) LinkOIDC(ctx context.Context, userID, issuer, subject string) er
|
||||
return nil
|
||||
}
|
||||
|
||||
// ListAdminIDs returns ids of all active admins (system notifications).
|
||||
func (s *Store) ListAdminIDs(ctx context.Context) ([]string, error) {
|
||||
cur, err := s.users().Find(ctx, bson.M{"roles.admin": true, "disabled": false})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("list admins: %w", err)
|
||||
}
|
||||
var users []domain.User
|
||||
if err := cur.All(ctx, &users); err != nil {
|
||||
return nil, fmt.Errorf("decode admins: %w", err)
|
||||
}
|
||||
ids := make([]string, len(users))
|
||||
for i, u := range users {
|
||||
ids[i] = u.ID
|
||||
}
|
||||
return ids, nil
|
||||
}
|
||||
|
||||
// TouchLastSeen updates lastSeenAt without bumping version (pure telemetry).
|
||||
func (s *Store) TouchLastSeen(ctx context.Context, userID string) error {
|
||||
_, err := s.users().UpdateOne(ctx, bson.M{"_id": userID},
|
||||
|
||||
Reference in New Issue
Block a user