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