Files
cgalo5758 a94ff08336 Add Discourse integration
Deliver forum posting entitlements through managed group membership with
identity linkage, periodic reconciliation, webhook handling, and an
operator mapping surface.

Include fake and live test environments, setup documentation,
migrations,
and end-to-end coverage.
2026-07-20 19:49:34 -07:00

144 lines
4.9 KiB
Go

package workflows
import (
"context"
"errors"
"fmt"
"log/slog"
"time"
"git.coopcloud.tech/wiki-cafe/member-console/internal/workflows/queues"
"go.temporal.io/api/serviceerror"
"go.temporal.io/sdk/client"
"go.temporal.io/sdk/temporal"
)
const (
// SyncScheduleID is the unique identifier for the Discourse group-sync
// schedule.
SyncScheduleID = "discourse-group-sync"
// SyncWorkflowID is the workflow ID used for sweep executions.
SyncWorkflowID = "discourse-group-sync-workflow"
// DefaultSyncInterval is the default sweep cadence. The sweep is the
// correctness safety net, not the latency path (webhooks and targeted
// reconciles are), so it runs at a low cadence with jitter to stay
// polite toward the site-wide shared admin rate bucket.
DefaultSyncInterval = 15 * time.Minute
)
// ScheduleConfig holds configuration for the group-sync schedule.
type ScheduleConfig struct {
Interval time.Duration
TriggerImmediately bool
}
// ScheduleManager manages the Discourse group-sync schedule. Its
// ensure/describe/create/update shape mirrors FedWiki's hardened
// ScheduleManager: Describe first, treat only NotFound as absence, and
// recover a lost create race with a single update.
type ScheduleManager struct {
client client.Client
logger *slog.Logger
}
// NewScheduleManager creates a ScheduleManager.
func NewScheduleManager(c client.Client, logger *slog.Logger) *ScheduleManager {
return &ScheduleManager{client: c, logger: logger}
}
// EnsureSyncSchedule creates or updates the group-sync schedule.
func (m *ScheduleManager) EnsureSyncSchedule(ctx context.Context, cfg ScheduleConfig) error {
if cfg.Interval == 0 {
cfg.Interval = DefaultSyncInterval
}
handle := m.client.ScheduleClient().GetHandle(ctx, SyncScheduleID)
_, err := handle.Describe(ctx)
switch {
case err == nil:
return m.updateSchedule(ctx, handle, cfg)
case isScheduleNotFound(err):
if createErr := m.createSchedule(ctx, cfg); createErr != nil {
if errors.Is(createErr, temporal.ErrScheduleAlreadyRunning) {
return m.updateSchedule(ctx, handle, cfg)
}
return fmt.Errorf("failed to create discourse sync schedule: %w", createErr)
}
return nil
default:
// A transient Describe failure is NOT "doesn't exist" — surface it
// rather than misrouting to Create.
return fmt.Errorf("failed to describe discourse sync schedule: %w", err)
}
}
func isScheduleNotFound(err error) bool {
var nf *serviceerror.NotFound
return errors.As(err, &nf)
}
// scheduleSpec builds the interval spec with jitter (a tenth of the
// interval) so multiple deployments against one forum don't synchronize
// their sweeps.
func scheduleSpec(cfg ScheduleConfig) client.ScheduleSpec {
return client.ScheduleSpec{
Intervals: []client.ScheduleIntervalSpec{{Every: cfg.Interval}},
Jitter: cfg.Interval / 10,
}
}
func scheduleAction() *client.ScheduleWorkflowAction {
return &client.ScheduleWorkflowAction{
ID: SyncWorkflowID,
Workflow: DiscourseGroupSyncWorkflow,
TaskQueue: queues.Main,
Args: []interface{}{SweepInput{}},
}
}
func (m *ScheduleManager) createSchedule(ctx context.Context, cfg ScheduleConfig) error {
m.logger.Info("creating discourse group-sync schedule",
slog.Duration("interval", cfg.Interval),
slog.Bool("triggerImmediately", cfg.TriggerImmediately))
spec := scheduleSpec(cfg)
if _, err := m.client.ScheduleClient().Create(ctx, client.ScheduleOptions{
ID: SyncScheduleID,
Spec: spec,
Action: scheduleAction(),
TriggerImmediately: cfg.TriggerImmediately,
}); err != nil {
return err
}
m.logger.Info("discourse group-sync schedule created", slog.String("scheduleID", SyncScheduleID))
return nil
}
func (m *ScheduleManager) updateSchedule(ctx context.Context, handle client.ScheduleHandle, cfg ScheduleConfig) error {
m.logger.Info("updating discourse group-sync schedule", slog.Duration("interval", cfg.Interval))
if err := handle.Update(ctx, client.ScheduleUpdateOptions{
DoUpdate: func(schedule client.ScheduleUpdateInput) (*client.ScheduleUpdate, error) {
spec := scheduleSpec(cfg)
schedule.Description.Schedule.Spec = &spec
schedule.Description.Schedule.Action = scheduleAction()
return &client.ScheduleUpdate{Schedule: &schedule.Description.Schedule}, nil
},
}); err != nil {
return fmt.Errorf("failed to update discourse sync schedule: %w", err)
}
m.logger.Info("discourse group-sync schedule updated")
return nil
}
// TriggerSyncNow triggers an immediate sweep (used by the operator surface
// and by tests).
func (m *ScheduleManager) TriggerSyncNow(ctx context.Context) error {
handle := m.client.ScheduleClient().GetHandle(ctx, SyncScheduleID)
if err := handle.Trigger(ctx, client.ScheduleTriggerOptions{}); err != nil {
return fmt.Errorf("failed to trigger discourse sync: %w", err)
}
m.logger.Info("discourse group-sync triggered", slog.String("scheduleID", SyncScheduleID))
return nil
}