Files
cgalo5758 88db730fcc Add dual licensing and SPDX headers
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.
2026-09-06 02:29:42 -05:00

223 lines
7.7 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 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
}