Files
cgalo5758 6f1effe58e Harden FedWiki sync schedule recovery
Only create schedules after Temporal NotFound and update after
concurrent
creation. Resolve the system workspace on every workflow run so
persisted
schedule arguments cannot retain stale IDs after database resets.
2026-07-11 16:08:52 -05:00

220 lines
7.6 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 FedWiki site sync schedule.
SyncScheduleID = "fedwiki-site-sync"
// SyncWorkflowID is the workflow ID used for sync workflow executions.
SyncWorkflowID = "fedwiki-site-sync-workflow"
// DefaultSyncInterval is the default interval between site syncs.
DefaultSyncInterval = 1 * time.Hour
)
// ScheduleConfig holds configuration for the FedWiki sync schedule.
type ScheduleConfig struct {
// Interval is how often to run the sync workflow.
Interval time.Duration
// TriggerImmediately runs the sync immediately when creating the schedule.
TriggerImmediately bool
}
// DefaultScheduleConfig returns a default schedule configuration.
func DefaultScheduleConfig() ScheduleConfig {
return ScheduleConfig{
Interval: DefaultSyncInterval,
TriggerImmediately: true,
}
}
// ScheduleManager manages Temporal schedules for FedWiki operations.
type ScheduleManager struct {
client client.Client
logger *slog.Logger
}
// NewScheduleManager creates a new ScheduleManager.
func NewScheduleManager(c client.Client, logger *slog.Logger) *ScheduleManager {
return &ScheduleManager{
client: c,
logger: logger,
}
}
// EnsureSyncSchedule creates or updates the FedWiki site sync schedule.
// It describes the existing schedule first: on success it updates it; a Temporal
// NotFound means the schedule is absent and it is created. Any other Describe
// error is surfaced to the caller rather than misread as absence (which would
// wrongly route to Create). A Create that loses a race against a concurrent
// creation falls back to a single Update so the schedule's args are still
// refreshed.
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) {
// Lost a create race — the schedule exists now; update it once.
return m.updateSchedule(ctx, handle, cfg)
}
return fmt.Errorf("failed to create sync schedule: %w", createErr)
}
return nil
default:
// Any other Describe failure (timeout, degraded Temporal, auth race) is
// NOT "doesn't exist" — surface it instead of masking it as create.
return fmt.Errorf("failed to describe sync schedule: %w", err)
}
}
// isScheduleNotFound reports whether err is Temporal's NotFound response. A
// schedule is backed by a workflow execution, so its "absent" signal arrives as
// *serviceerror.NotFound — distinct from a transient describe failure.
func isScheduleNotFound(err error) bool {
var nf *serviceerror.NotFound
return errors.As(err, &nf)
}
// createSchedule registers the sync schedule. Args carry an empty
// SyncFedWikiSitesWorkflowInput — the System workspace is resolved per execution
// by the workflow, not baked into schedule args.
func (m *ScheduleManager) createSchedule(ctx context.Context, cfg ScheduleConfig) error {
m.logger.Info("creating FedWiki sync schedule",
slog.Duration("interval", cfg.Interval),
slog.Bool("triggerImmediately", cfg.TriggerImmediately))
if _, err := m.client.ScheduleClient().Create(ctx, client.ScheduleOptions{
ID: SyncScheduleID,
Spec: client.ScheduleSpec{
Intervals: []client.ScheduleIntervalSpec{
{Every: cfg.Interval},
},
},
Action: &client.ScheduleWorkflowAction{
ID: SyncWorkflowID,
Workflow: SyncFedWikiSitesWorkflow,
TaskQueue: queues.Main,
Args: []interface{}{SyncFedWikiSitesWorkflowInput{}},
},
TriggerImmediately: cfg.TriggerImmediately,
}); err != nil {
return err
}
m.logger.Info("FedWiki sync schedule created successfully",
slog.String("scheduleID", SyncScheduleID))
return nil
}
// updateSchedule rewrites the existing sync schedule's spec and action. Args
// carry an empty SyncFedWikiSitesWorkflowInput, which strips any stale
// DefaultWorkspaceID baked into a schedule created under the previous code.
func (m *ScheduleManager) updateSchedule(ctx context.Context, handle client.ScheduleHandle, cfg ScheduleConfig) error {
m.logger.Info("updating existing FedWiki sync schedule",
slog.Duration("interval", cfg.Interval))
if err := handle.Update(ctx, client.ScheduleUpdateOptions{
DoUpdate: func(schedule client.ScheduleUpdateInput) (*client.ScheduleUpdate, error) {
schedule.Description.Schedule.Spec = &client.ScheduleSpec{
Intervals: []client.ScheduleIntervalSpec{
{Every: cfg.Interval},
},
}
schedule.Description.Schedule.Action = &client.ScheduleWorkflowAction{
ID: SyncWorkflowID,
Workflow: SyncFedWikiSitesWorkflow,
TaskQueue: queues.Main,
Args: []interface{}{SyncFedWikiSitesWorkflowInput{}},
}
return &client.ScheduleUpdate{
Schedule: &schedule.Description.Schedule,
}, nil
},
}); err != nil {
return fmt.Errorf("failed to update sync schedule: %w", err)
}
m.logger.Info("FedWiki sync schedule updated successfully")
return nil
}
// DeleteSyncSchedule deletes the FedWiki site sync schedule.
func (m *ScheduleManager) DeleteSyncSchedule(ctx context.Context) error {
handle := m.client.ScheduleClient().GetHandle(ctx, SyncScheduleID)
err := handle.Delete(ctx)
if err != nil {
return fmt.Errorf("failed to delete sync schedule: %w", err)
}
m.logger.Info("FedWiki sync schedule deleted", slog.String("scheduleID", SyncScheduleID))
return nil
}
// PauseSyncSchedule pauses the FedWiki site sync schedule.
func (m *ScheduleManager) PauseSyncSchedule(ctx context.Context, note string) error {
handle := m.client.ScheduleClient().GetHandle(ctx, SyncScheduleID)
err := handle.Pause(ctx, client.SchedulePauseOptions{
Note: note,
})
if err != nil {
return fmt.Errorf("failed to pause sync schedule: %w", err)
}
m.logger.Info("FedWiki sync schedule paused",
slog.String("scheduleID", SyncScheduleID),
slog.String("note", note))
return nil
}
// UnpauseSyncSchedule unpauses (resumes) the FedWiki site sync schedule.
func (m *ScheduleManager) UnpauseSyncSchedule(ctx context.Context, note string) error {
handle := m.client.ScheduleClient().GetHandle(ctx, SyncScheduleID)
err := handle.Unpause(ctx, client.ScheduleUnpauseOptions{
Note: note,
})
if err != nil {
return fmt.Errorf("failed to unpause sync schedule: %w", err)
}
m.logger.Info("FedWiki sync schedule unpaused",
slog.String("scheduleID", SyncScheduleID),
slog.String("note", note))
return nil
}
// TriggerSyncNow triggers an immediate execution of the sync workflow.
func (m *ScheduleManager) TriggerSyncNow(ctx context.Context) error {
handle := m.client.ScheduleClient().GetHandle(ctx, SyncScheduleID)
err := handle.Trigger(ctx, client.ScheduleTriggerOptions{})
if err != nil {
return fmt.Errorf("failed to trigger sync schedule: %w", err)
}
m.logger.Info("FedWiki sync triggered", slog.String("scheduleID", SyncScheduleID))
return nil
}
// GetSyncScheduleInfo returns information about the sync schedule.
func (m *ScheduleManager) GetSyncScheduleInfo(ctx context.Context) (*client.ScheduleDescription, error) {
handle := m.client.ScheduleClient().GetHandle(ctx, SyncScheduleID)
desc, err := handle.Describe(ctx)
if err != nil {
return nil, fmt.Errorf("failed to describe sync schedule: %w", err)
}
return desc, nil
}