Introduce a commercial license option alongside AGPL-3.0-only, require a CLA for contributors, and document the terms in COMMERCIAL.md and NOTICE. Add a script to stamp SPDX headers on Go files and apply it across the tree.
147 lines
5.0 KiB
Go
147 lines
5.0 KiB
Go
// SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Commercial
|
|
// SPDX-FileCopyrightText: 2025-2026 Christian Galo
|
|
|
|
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
|
|
}
|