e8beb7d4f3
- Navigation: new per-user setting (default "side") renders the page links as a left sidebar while keeping theme/bell/avatar/logout in the top-right; falls back to a top bar under 760px or when set to "top". Profile selector added. - Group conversations: record a creatorId; the creator gets a "Manage members" button to add/remove participants (new GET/POST/DELETE members endpoints, creator-only). Existing groups backfilled. - Customer ticketing: admin can map each assigned consultant to their username in that customer's ticketing system. Per-customer mapping takes precedence over the consultant's global ticketingIdentities during sync. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
466 lines
13 KiB
Go
466 lines
13 KiB
Go
package sync
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"fmt"
|
|
"log/slog"
|
|
"net/http"
|
|
"sync"
|
|
"time"
|
|
|
|
"go.mongodb.org/mongo-driver/v2/bson"
|
|
|
|
"bountyboard/internal/domain"
|
|
"bountyboard/internal/files"
|
|
"bountyboard/internal/metrics"
|
|
"bountyboard/internal/store"
|
|
)
|
|
|
|
// overlapWindow re-fetches a little history each poll so clock skew or slow
|
|
// writes can't drop updates (§5.3).
|
|
const overlapWindow = 5 * time.Minute
|
|
|
|
// Publisher pushes live updates (wired to the WS hub when it exists).
|
|
type Publisher func(channel, event string, payload any)
|
|
|
|
// Manager runs one polling worker per syncable customer and reconciles the
|
|
// worker set against the customers collection.
|
|
type Manager struct {
|
|
st *store.Store
|
|
files *files.Store
|
|
log *slog.Logger
|
|
reg *metrics.Registry
|
|
creds func(*domain.Customer) (Credentials, error)
|
|
publish Publisher
|
|
|
|
mu sync.Mutex
|
|
workers map[string]*worker
|
|
}
|
|
|
|
type worker struct {
|
|
customerID string
|
|
interval time.Duration
|
|
trigger chan struct{}
|
|
cancel context.CancelFunc
|
|
|
|
mu sync.Mutex
|
|
name string
|
|
lastStatus string
|
|
lastError string
|
|
lastRunAt time.Time
|
|
}
|
|
|
|
// WorkerStatus feeds the admin service-status panel (§11.17).
|
|
type WorkerStatus struct {
|
|
CustomerID string `json:"customerId"`
|
|
CustomerName string `json:"customerName"`
|
|
PollIntervalSec int `json:"pollIntervalSec"`
|
|
LastStatus string `json:"lastStatus"`
|
|
LastError string `json:"lastError,omitempty"`
|
|
LastRunAt time.Time `json:"lastRunAt"`
|
|
}
|
|
|
|
func NewManager(st *store.Store, fs *files.Store, log *slog.Logger, reg *metrics.Registry,
|
|
creds func(*domain.Customer) (Credentials, error), publish Publisher) *Manager {
|
|
if publish == nil {
|
|
publish = func(string, string, any) {}
|
|
}
|
|
return &Manager{
|
|
st: st, files: fs, log: log, reg: reg, creds: creds, publish: publish,
|
|
workers: map[string]*worker{},
|
|
}
|
|
}
|
|
|
|
// Run reconciles workers until ctx is done, then stops them all (graceful
|
|
// shutdown flushes by waiting for the in-flight tick via context).
|
|
func (m *Manager) Run(ctx context.Context) {
|
|
ticker := time.NewTicker(30 * time.Second)
|
|
defer ticker.Stop()
|
|
m.reconcile(ctx)
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
m.mu.Lock()
|
|
for _, w := range m.workers {
|
|
w.cancel()
|
|
}
|
|
m.mu.Unlock()
|
|
return
|
|
case <-ticker.C:
|
|
m.reconcile(ctx)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (m *Manager) reconcile(ctx context.Context) {
|
|
customers, err := m.st.ListCustomers(ctx, "", false)
|
|
if err != nil {
|
|
m.log.Error("sync reconcile: list customers", "err", err)
|
|
return
|
|
}
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
|
|
want := map[string]*domain.Customer{}
|
|
for i := range customers {
|
|
c := &customers[i]
|
|
if c.Ticketing.Type.Valid() && len(c.ConsultantIDs) > 0 {
|
|
want[c.ID] = c
|
|
}
|
|
}
|
|
// stop removed/archived customers
|
|
for id, w := range m.workers {
|
|
if _, ok := want[id]; !ok {
|
|
w.cancel()
|
|
delete(m.workers, id)
|
|
}
|
|
}
|
|
// start new / restart changed-interval workers
|
|
for id, c := range want {
|
|
interval := time.Duration(c.Ticketing.PollIntervalSec) * time.Second
|
|
if interval < 10*time.Second {
|
|
interval = time.Minute
|
|
}
|
|
if w, ok := m.workers[id]; ok {
|
|
if w.interval == interval {
|
|
w.mu.Lock()
|
|
w.name = c.Name
|
|
w.mu.Unlock()
|
|
continue
|
|
}
|
|
w.cancel()
|
|
delete(m.workers, id)
|
|
}
|
|
wctx, cancel := context.WithCancel(ctx)
|
|
w := &worker{
|
|
customerID: id, interval: interval, name: c.Name,
|
|
trigger: make(chan struct{}, 1), cancel: cancel,
|
|
}
|
|
m.workers[id] = w
|
|
go m.runWorker(wctx, w)
|
|
}
|
|
m.reg.Set("sync_workers", int64(len(m.workers)))
|
|
}
|
|
|
|
// SyncNow triggers an immediate poll for the customer.
|
|
func (m *Manager) SyncNow(customerID string) {
|
|
m.mu.Lock()
|
|
w, ok := m.workers[customerID]
|
|
m.mu.Unlock()
|
|
if ok {
|
|
select {
|
|
case w.trigger <- struct{}{}:
|
|
default: // already queued
|
|
}
|
|
}
|
|
}
|
|
|
|
// Statuses implements the admin panel provider.
|
|
func (m *Manager) Statuses() any {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
out := []WorkerStatus{}
|
|
for _, w := range m.workers {
|
|
w.mu.Lock()
|
|
out = append(out, WorkerStatus{
|
|
CustomerID: w.customerID,
|
|
CustomerName: w.name,
|
|
PollIntervalSec: int(w.interval.Seconds()),
|
|
LastStatus: w.lastStatus,
|
|
LastError: w.lastError,
|
|
LastRunAt: w.lastRunAt,
|
|
})
|
|
w.mu.Unlock()
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (m *Manager) runWorker(ctx context.Context, w *worker) {
|
|
// first run promptly, then on the interval
|
|
timer := time.NewTimer(2 * time.Second)
|
|
defer timer.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-timer.C:
|
|
case <-w.trigger:
|
|
}
|
|
m.syncOnce(ctx, w)
|
|
timer.Reset(w.interval)
|
|
}
|
|
}
|
|
|
|
func (m *Manager) syncOnce(ctx context.Context, w *worker) {
|
|
defer func() {
|
|
if rec := recover(); rec != nil {
|
|
m.log.Error("sync worker panic", "customerId", w.customerID, "panic", fmt.Sprint(rec))
|
|
}
|
|
}()
|
|
ctx, cancel := context.WithTimeout(ctx, 2*time.Minute)
|
|
defer cancel()
|
|
|
|
c, err := m.st.CustomerByID(ctx, w.customerID)
|
|
if err != nil {
|
|
m.recordRun(ctx, w, nil, err)
|
|
return
|
|
}
|
|
err = m.SyncCustomer(ctx, c)
|
|
m.recordRun(ctx, w, c, err)
|
|
}
|
|
|
|
func (m *Manager) recordRun(ctx context.Context, w *worker, c *domain.Customer, err error) {
|
|
now := time.Now().UTC()
|
|
status, msg := "ok", ""
|
|
if err != nil {
|
|
status, msg = "error", err.Error()
|
|
m.log.Error("sync run failed", "customerId", w.customerID, "err", err)
|
|
m.reg.Inc("sync_runs_failed_total", 1)
|
|
// surface as admin notification (§12) — best effort, throttled by
|
|
// only notifying on status flips
|
|
if c != nil && c.Ticketing.LastSyncStatus != "error" {
|
|
m.notifyAdmins(ctx, "sync_error", "Sync failed: "+c.Name, msg, "/admin")
|
|
}
|
|
} else {
|
|
m.reg.Inc("sync_runs_total", 1)
|
|
}
|
|
w.mu.Lock()
|
|
w.lastStatus, w.lastError, w.lastRunAt = status, msg, now
|
|
w.mu.Unlock()
|
|
if setErr := m.st.SetSyncStatus(ctx, w.customerID, now, status, msg); setErr != nil {
|
|
m.log.Warn("set sync status", "err", setErr)
|
|
}
|
|
}
|
|
|
|
// SyncCustomer performs one full sync pass for a customer; exported for
|
|
// tests and the sync-now path.
|
|
func (m *Manager) SyncCustomer(ctx context.Context, c *domain.Customer) error {
|
|
creds, err := m.creds(c)
|
|
if err != nil {
|
|
return fmt.Errorf("decrypt credentials: %w", err)
|
|
}
|
|
conn, err := NewConnector(c.Ticketing.Type, c.Ticketing.BaseURL, c.Ticketing.ProjectKey, creds)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
since := c.Ticketing.LastSyncAt.Add(-overlapWindow)
|
|
if c.Ticketing.LastSyncAt.IsZero() {
|
|
since = time.Unix(0, 0)
|
|
}
|
|
|
|
var firstErr error
|
|
for _, consultantID := range c.ConsultantIDs {
|
|
consultant, err := m.st.UserByID(ctx, consultantID)
|
|
if err != nil {
|
|
continue // consultant deleted; reconcile will catch up
|
|
}
|
|
// Per-customer mapping (set by the admin on the customer) wins over
|
|
// the consultant's global identity; fall back to the global one.
|
|
identity := c.Ticketing.ConsultantIdentities[consultantID]
|
|
if identity == "" {
|
|
identity = ticketingIdentity(consultant, c.Ticketing.Type)
|
|
}
|
|
if identity == "" {
|
|
continue // §5.3: consultant has no linked account in this system
|
|
}
|
|
|
|
tickets, err := conn.FetchUpdated(ctx, identity, since)
|
|
if err != nil {
|
|
if firstErr == nil {
|
|
firstErr = fmt.Errorf("fetch for %s: %w", consultant.Email, err)
|
|
}
|
|
continue
|
|
}
|
|
for _, ticket := range tickets {
|
|
if err := m.upsertTicket(ctx, c, consultant, conn, ticket); err != nil {
|
|
m.log.Error("upsert ticket", "key", ticket.Key, "err", err)
|
|
if firstErr == nil {
|
|
firstErr = err
|
|
}
|
|
}
|
|
}
|
|
|
|
if err := m.flagOrphans(ctx, c, consultantID, conn, identity); err != nil {
|
|
m.log.Warn("orphan detection", "customerId", c.ID, "err", err)
|
|
}
|
|
}
|
|
return firstErr
|
|
}
|
|
|
|
func ticketingIdentity(u *domain.User, t domain.TicketingType) string {
|
|
if t == domain.TicketingDemo {
|
|
return u.Email // demo accepts any identity
|
|
}
|
|
// The driver decodes the nested document as bson.D when the field type
|
|
// is `any`; JSON-built values arrive as map[string]any. Handle both.
|
|
switch ids := u.Extra["ticketingIdentities"].(type) {
|
|
case map[string]any:
|
|
v, _ := ids[string(t)].(string)
|
|
return v
|
|
case bson.M:
|
|
v, _ := ids[string(t)].(string)
|
|
return v
|
|
case bson.D:
|
|
for _, e := range ids {
|
|
if e.Key == string(t) {
|
|
v, _ := e.Value.(string)
|
|
return v
|
|
}
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func contentHash(t Ticket) string {
|
|
h := sha256.New()
|
|
fmt.Fprintf(h, "%s|%s|%s|%d", t.Title, t.Description, t.Type, len(t.Attachments))
|
|
for _, a := range t.Attachments {
|
|
fmt.Fprintf(h, "|%s", a.Name)
|
|
}
|
|
return hex.EncodeToString(h.Sum(nil))
|
|
}
|
|
|
|
func (m *Manager) upsertTicket(ctx context.Context, c *domain.Customer,
|
|
consultant *domain.User, conn Connector, ticket Ticket) error {
|
|
|
|
task := &domain.Task{
|
|
CustomerID: c.ID,
|
|
ConsultantID: consultant.ID,
|
|
Title: ticket.Title,
|
|
Description: ticket.Description,
|
|
Budget: c.DefaultBudget, // §11.20 pre-fill
|
|
EffortCoefficient: 1,
|
|
External: &domain.External{
|
|
System: string(c.Ticketing.Type),
|
|
Key: ticket.Key,
|
|
URL: ticket.URL,
|
|
Type: ticket.Type,
|
|
Raw: ticket.Raw,
|
|
ContentHash: contentHash(ticket),
|
|
},
|
|
}
|
|
task.Bounty = domain.ComputeBounty(task.EffortCoefficient, task.Budget)
|
|
for _, a := range ticket.Attachments {
|
|
task.Attachments = append(task.Attachments, domain.TaskAttachment{
|
|
Name: a.Name, URL: a.URL, MimeType: a.MimeType,
|
|
})
|
|
}
|
|
|
|
res, err := m.st.UpsertImportedTask(ctx, task)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
switch {
|
|
case res.Created:
|
|
m.reg.Inc("sync_tickets_imported_total", 1)
|
|
m.cacheAttachments(ctx, conn, res.Task)
|
|
m.notify(ctx, consultant.ID, "ticket_imported",
|
|
"New ticket imported: "+ticket.Key,
|
|
ticket.Title, "/tasks/"+res.Task.ID)
|
|
m.publish("board", "task.imported", map[string]any{"taskId": res.Task.ID, "customerId": c.ID})
|
|
case res.UpstreamChanged:
|
|
m.notify(ctx, consultant.ID, "upstream_changed",
|
|
"Upstream ticket changed: "+ticket.Key,
|
|
"The source ticket changed after atomization began.", "/tasks/"+res.Task.ID)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// cacheAttachments downloads ticket attachments into GridFS so the atomizer
|
|
// can fetch them from the app (§5.3); failures are logged, not fatal.
|
|
func (m *Manager) cacheAttachments(ctx context.Context, conn Connector, t *domain.Task) {
|
|
changed := false
|
|
for i, a := range t.Attachments {
|
|
if a.FileID != "" || a.URL == "" {
|
|
continue
|
|
}
|
|
fileID, err := m.downloadAttachment(ctx, conn, t.ID, a)
|
|
if err != nil {
|
|
m.log.Warn("cache attachment", "task", t.ID, "name", a.Name, "err", err)
|
|
continue
|
|
}
|
|
t.Attachments[i].FileID = fileID
|
|
changed = true
|
|
}
|
|
if changed {
|
|
if err := m.st.SetTaskAttachments(ctx, t.ID, t.Attachments); err != nil {
|
|
m.log.Warn("persist attachment cache", "task", t.ID, "err", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (m *Manager) downloadAttachment(ctx context.Context, conn Connector, taskID string, a domain.TaskAttachment) (string, error) {
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, a.URL, nil)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
conn.Authorize(req)
|
|
resp, err := httpClient.Do(req)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode != http.StatusOK {
|
|
return "", fmt.Errorf("download %s: status %d", a.URL, resp.StatusCode)
|
|
}
|
|
meta, err := m.files.Save(ctx, files.ScopeTask, "system:"+taskID, a.Name, a.MimeType, resp.Body)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
return meta.ID, nil
|
|
}
|
|
|
|
func (m *Manager) flagOrphans(ctx context.Context, c *domain.Customer,
|
|
consultantID string, conn Connector, identity string) error {
|
|
|
|
keys, err := conn.ListAssignedKeys(ctx, identity)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
current := make(map[string]bool, len(keys))
|
|
for _, k := range keys {
|
|
current[k] = true
|
|
}
|
|
tasks, err := m.st.ListImportedTasksForOrphanCheck(ctx, c.ID, string(c.Ticketing.Type), consultantID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, t := range tasks {
|
|
if t.External == nil || current[t.External.Key] || t.External.Orphaned {
|
|
continue
|
|
}
|
|
if err := m.st.MarkTaskOrphaned(ctx, t.ID); err != nil {
|
|
continue // raced with another worker — fine
|
|
}
|
|
m.reg.Inc("sync_tickets_orphaned_total", 1)
|
|
m.notify(ctx, consultantID, "ticket_orphaned",
|
|
"Ticket disappeared upstream: "+t.External.Key,
|
|
t.Title+" — decide whether to archive it.", "/tasks/"+t.ID)
|
|
m.publish("board", "task.orphaned", map[string]any{"taskId": t.ID, "customerId": c.ID})
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (m *Manager) notify(ctx context.Context, userID, kind, title, body, link string) {
|
|
n := &store.Notification{UserID: userID, Kind: kind, Title: title, Body: body, Link: link}
|
|
if err := m.st.InsertNotification(ctx, n); err != nil {
|
|
m.log.Warn("insert notification", "err", err)
|
|
return
|
|
}
|
|
m.publish("notifications", "notification", n)
|
|
}
|
|
|
|
func (m *Manager) notifyAdmins(ctx context.Context, kind, title, body, link string) {
|
|
admins, err := m.st.ListAdminIDs(ctx)
|
|
if err != nil {
|
|
m.log.Warn("list admins", "err", err)
|
|
return
|
|
}
|
|
for _, id := range admins {
|
|
m.notify(ctx, id, kind, title, body, link)
|
|
}
|
|
}
|