diff --git a/backend/cache/local.go b/backend/cache/local.go index 0968030..85d8d7e 100644 --- a/backend/cache/local.go +++ b/backend/cache/local.go @@ -51,10 +51,11 @@ var CampaignEventPriority = map[string]int{ data.EVENT_CAMPAIGN_RECIPIENT_MESSAGE_SENT: 20, data.EVENT_CAMPAIGN_RECIPIENT_SCHEDULED: 10, // campaign events - data.EVENT_CAMPAIGN_CLOSED: 30, - data.EVENT_CAMPAIGN_ACTIVE: 20, - data.EVENT_CAMPAIGN_SELF_MANAGED: 20, - data.EVENT_CAMPAIGN_SCHEDULED: 10, + data.EVENT_CAMPAIGN_CLOSED: 30, + data.EVENT_CAMPAIGN_ACTIVE: 20, + data.EVENT_CAMPAIGN_SELF_MANAGED: 20, + data.EVENT_CAMPAIGN_SCHEDULED: 10, + data.EVENT_CAMPAIGN_PENDING_SCHEDULE: 5, } // IsMoreNotableCampaignRecipientEvent returns true if newEvent is more notable than currentEvent diff --git a/backend/data/events.go b/backend/data/events.go index b9e1e50..8bd443a 100644 --- a/backend/data/events.go +++ b/backend/data/events.go @@ -1,10 +1,11 @@ package data const ( - EVENT_CAMPAIGN_SCHEDULED = "campaign_scheduled" - EVENT_CAMPAIGN_ACTIVE = "campaign_active" - EVENT_CAMPAIGN_SELF_MANAGED = "campaign_self_managed" - EVENT_CAMPAIGN_CLOSED = "campaign_closed" + EVENT_CAMPAIGN_SCHEDULED = "campaign_scheduled" + EVENT_CAMPAIGN_ACTIVE = "campaign_active" + EVENT_CAMPAIGN_SELF_MANAGED = "campaign_self_managed" + EVENT_CAMPAIGN_PENDING_SCHEDULE = "campaign_pending_schedule" + EVENT_CAMPAIGN_CLOSED = "campaign_closed" EVENT_CAMPAIGN_RECIPIENT_SCHEDULED = "campaign_recipient_scheduled" EVENT_CAMPAIGN_RECIPIENT_MESSAGE_SENT = "campaign_recipient_message_sent" @@ -26,6 +27,7 @@ var Events = []string{ EVENT_CAMPAIGN_SCHEDULED, EVENT_CAMPAIGN_ACTIVE, EVENT_CAMPAIGN_SELF_MANAGED, + EVENT_CAMPAIGN_PENDING_SCHEDULE, EVENT_CAMPAIGN_CLOSED, // campaign recipient events EVENT_CAMPAIGN_RECIPIENT_SCHEDULED, diff --git a/backend/database/campaign.go b/backend/database/campaign.go index 7b48f36..f99cf25 100644 --- a/backend/database/campaign.go +++ b/backend/database/campaign.go @@ -25,6 +25,15 @@ type Campaign struct { SortOrder string `gorm:";"` // 'asc,desc,random' SendStartAt *time.Time `gorm:"index;"` SendEndAt *time.Time `gorm:"index;"` + // ScheduleAt is set when the campaign uses late-scheduling. + // the task runner will call schedule() when now >= ScheduleAt. + // null means the campaign was scheduled immediately at creation. + ScheduleAt *time.Time `gorm:"index;"` + // JitterMin and JitterMax are persisted only when ScheduleAt is set (late-scheduling). + // They are cleared after the campaign is scheduled by the task runner. + // For immediately-scheduled campaigns these columns remain null. + JitterMin *int `gorm:""` + JitterMax *int `gorm:""` // ConstraintWeekDays is a binary format. // 0b00000001 = 1 = sunday diff --git a/backend/model/campaign.go b/backend/model/campaign.go index 43dd35d..2ab4774 100644 --- a/backend/model/campaign.go +++ b/backend/model/campaign.go @@ -28,11 +28,14 @@ type Campaign struct { SortOrder nullable.Nullable[vo.CampaignSendingOrder] `json:"sortOrder"` SendStartAt nullable.Nullable[time.Time] `json:"sendStartAt"` SendEndAt nullable.Nullable[time.Time] `json:"sendEndAt"` + ScheduleAt nullable.Nullable[time.Time] `json:"scheduleAt"` ConstraintWeekDays nullable.Nullable[vo.CampaignWeekDays] `json:"constraintWeekDays"` ConstraintStartTime nullable.Nullable[vo.CampaignTimeConstraint] `json:"constraintStartTime"` ConstraintEndTime nullable.Nullable[vo.CampaignTimeConstraint] `json:"constraintEndTime"` - // jitter is used only during scheduling, not persisted to database + // jitter is persisted to the database only when ScheduleAt is set (late-scheduling), + // so the task runner can apply it when schedule() is called hours later. + // For immediately-scheduled campaigns jitter lives in memory only and is never written to the DB. JitterMin nullable.Nullable[int] `json:"jitterMin,omitempty"` JitterMax nullable.Nullable[int] `json:"jitterMax,omitempty"` @@ -410,6 +413,12 @@ func (c *Campaign) ToDBMap() map[string]any { m["send_end_at"] = utils.RFC3339UTC(v) } } + if c.ScheduleAt.IsSpecified() { + m["schedule_at"] = nil + if v, err := c.ScheduleAt.Get(); err == nil { + m["schedule_at"] = utils.RFC3339UTC(v) + } + } if c.CloseAt.IsSpecified() { m["close_at"] = nil if v, err := c.CloseAt.Get(); err == nil { @@ -504,6 +513,18 @@ func (c *Campaign) ToDBMap() map[string]any { if v, err := c.NotableEventID.Get(); err == nil { m["notable_event_id"] = v.String() } + if c.JitterMin.IsSpecified() { + m["jitter_min"] = nil + if v, err := c.JitterMin.Get(); err == nil { + m["jitter_min"] = v + } + } + if c.JitterMax.IsSpecified() { + m["jitter_max"] = nil + if v, err := c.JitterMax.Get(); err == nil { + m["jitter_max"] = v + } + } return m } diff --git a/backend/repository/campaign.go b/backend/repository/campaign.go index c08b9e6..2e893d7 100644 --- a/backend/repository/campaign.go +++ b/backend/repository/campaign.go @@ -1366,6 +1366,41 @@ func (r *Campaign) GetReadyToAnonymize( return result, nil } +// GetReadyToLateSchedule returns campaigns where schedule_at <= now and +// the campaign has not yet been scheduled (notable event is still pending_schedule). +func (r *Campaign) GetReadyToLateSchedule( + ctx context.Context, + options *CampaignOption, +) (*model.Result[model.Campaign], error) { + result := model.NewEmptyResult[model.Campaign]() + db := r.load(r.DB, options) + db, err := useQuery(db, database.CAMPAIGN_TABLE, options.QueryArgs) + if err != nil { + return result, errs.Wrap(err) + } + pendingEventID := cache.EventIDByName[data.EVENT_CAMPAIGN_PENDING_SCHEDULE] + var dbCampaigns []database.Campaign + res := db. + Where("schedule_at <= ? AND notable_event_id = ?", utils.NowRFC3339UTC(), pendingEventID). + Find(&dbCampaigns) + if res.Error != nil { + return result, res.Error + } + hasNextPage, err := useHasNextPage(db, database.CAMPAIGN_TABLE, options.QueryArgs) + if err != nil { + return result, errs.Wrap(err) + } + result.HasNextPage = hasNextPage + for _, dbCampaign := range dbCampaigns { + campaign, err := ToCampaign(&dbCampaign) + if err != nil { + return nil, errs.Wrap(err) + } + result.Rows = append(result.Rows, campaign) + } + return result, nil +} + // SaveEvent saves a campaign event func (r *Campaign) SaveEvent( ctx context.Context, @@ -1793,6 +1828,24 @@ func ToCampaign(row *database.Campaign) (*model.Campaign, error) { } else { sendEndAt.SetNull() } + var scheduleAt nullable.Nullable[time.Time] + if row.ScheduleAt != nil { + scheduleAt = nullable.NewNullableWithValue(*row.ScheduleAt) + } else { + scheduleAt.SetNull() + } + var jitterMin nullable.Nullable[int] + if row.JitterMin != nil { + jitterMin = nullable.NewNullableWithValue(*row.JitterMin) + } else { + jitterMin.SetNull() + } + var jitterMax nullable.Nullable[int] + if row.JitterMax != nil { + jitterMax = nullable.NewNullableWithValue(*row.JitterMax) + } else { + jitterMax.SetNull() + } saveSubmittedData := nullable.NewNullableWithValue(row.SaveSubmittedData) saveBrowserMetadata := nullable.NewNullableWithValue(row.SaveBrowserMetadata) isAnonymous := nullable.NewNullableWithValue(row.IsAnonymous) @@ -1927,6 +1980,9 @@ func ToCampaign(row *database.Campaign) (*model.Campaign, error) { SortOrder: sortOrder, SendStartAt: sendStartAt, SendEndAt: sendEndAt, + ScheduleAt: scheduleAt, + JitterMin: jitterMin, + JitterMax: jitterMax, ConstraintWeekDays: constraintWeekDays, ConstraintStartTime: constraintStartTime, ConstraintEndTime: constraintEndTime, diff --git a/backend/service/campaign.go b/backend/service/campaign.go index 6a71b05..5ec6e1b 100644 --- a/backend/service/campaign.go +++ b/backend/service/campaign.go @@ -118,6 +118,33 @@ func (c *Campaign) Create( return nil, errs.Wrap(err) } } + // late-schedule (indicated by a non-null schedule_at) is not compatible with self-managed campaigns + if campaign.ScheduleAt.IsSpecified() && !campaign.ScheduleAt.IsNull() { + if campaign.IsSelfManaged() { + return nil, validate.WrapErrorWithField( + errors.New("late scheduling is not available for self-managed campaigns"), + "scheduleAt", + ) + } + // scheduleAt must be in the future + scheduleAt := campaign.ScheduleAt.MustGet() + if !scheduleAt.After(time.Now().UTC()) { + return nil, validate.WrapErrorWithField( + errors.New("schedule time must be in the future"), + "scheduleAt", + ) + } + // send_start_at must be more than 24h after schedule_at + if campaign.SendStartAt.IsSpecified() && !campaign.SendStartAt.IsNull() { + sendStartAt := campaign.SendStartAt.MustGet() + if sendStartAt.Sub(scheduleAt) < 24*time.Hour { + return nil, validate.WrapErrorWithField( + errors.New("send start must be at least 24 hours after the schedule-at time"), + "scheduleAt", + ) + } + } + } // validate if err := campaign.Validate(); err != nil { return nil, errs.Wrap(err) @@ -213,14 +240,26 @@ func (c *Campaign) Create( c.Logger.Errorw("failed to get campaign by id", "error", err) return nil, errs.Wrap(err) } - // preserve jitter values from original campaign (not persisted to db) - createdCampaign.JitterMin = campaign.JitterMin - createdCampaign.JitterMax = campaign.JitterMax - err = c.schedule(ctx, session, createdCampaign) - if err != nil { - c.Logger.Errorw("failed to schedule campaign", "error", err) - // TODO we should delete the campaign as it was not scheduled - return nil, errs.Wrap(err) + if createdCampaign.ScheduleAt.IsSpecified() && !createdCampaign.ScheduleAt.IsNull() { + // Late-scheduling: persist jitter to the DB so the task runner can apply it + // when schedule() is called hours later. The in-memory copy on createdCampaign + // is set too so setMostNotableCampaignEvent's UpdateByID writes the columns. + createdCampaign.JitterMin = campaign.JitterMin + createdCampaign.JitterMax = campaign.JitterMax + err = c.setMostNotableCampaignEvent(ctx, createdCampaign, data.EVENT_CAMPAIGN_PENDING_SCHEDULE) + if err != nil { + return nil, errs.Wrap(err) + } + } else { + // Immediate scheduling: jitter lives in memory only for this call, never written to DB. + createdCampaign.JitterMin = campaign.JitterMin + createdCampaign.JitterMax = campaign.JitterMax + err = c.schedule(ctx, session, createdCampaign) + if err != nil { + c.Logger.Errorw("failed to schedule campaign", "error", err) + // TODO we should delete the campaign as it was not scheduled + return nil, errs.Wrap(err) + } } ae.Details["id"] = id.String() c.AuditLogAuthorized(ae) @@ -1482,6 +1521,15 @@ func (c *Campaign) UpdateByID( if v, err := incoming.AnonymizedAt.Get(); err == nil { current.AnonymizedAt.Set(v.Truncate(time.Minute)) } + if v, err := incoming.ScheduleAt.Get(); err == nil { + current.ScheduleAt.Set(v) + } else if incoming.ScheduleAt.IsSpecified() { + // incoming was explicitly null — clear the scheduled time, reverting to immediate scheduling + current.ScheduleAt.SetNull() + // also clear any persisted jitter — it was stored for late-scheduling and is no longer needed + current.JitterMin.SetNull() + current.JitterMax.SetNull() + } if v, err := incoming.RecipientGroupIDs.Get(); err == nil { current.RecipientGroupIDs.Set(v) } @@ -1565,6 +1613,47 @@ func (c *Campaign) UpdateByID( current.EvasionPageID.Set(incoming.EvasionPageID.MustGet()) } } + // late-schedule (indicated by a non-null schedule_at) is not compatible with self-managed campaigns + if current.ScheduleAt.IsSpecified() && !current.ScheduleAt.IsNull() { + if current.IsSelfManaged() { + return validate.WrapErrorWithField( + errors.New("late scheduling is not available for self-managed campaigns"), + "scheduleAt", + ) + } + // reject scheduleAt if the campaign has already moved past pending_schedule — + // at that point recipients have been resolved and scheduling is done; there is + // nothing meaningful for a new scheduleAt to do and it would revert the campaign + // back to pending_schedule state unexpectedly. + if currentNotableEventID, err := current.NotableEventID.Get(); err == nil { + if currentEventName, ok := cache.EventNameByID[currentNotableEventID.String()]; ok { + if cache.IsMoreNotableCampaignRecipientEvent(currentEventName, data.EVENT_CAMPAIGN_PENDING_SCHEDULE) { + return validate.WrapErrorWithField( + errors.New("scheduleAt cannot be set on a campaign that has already been scheduled"), + "scheduleAt", + ) + } + } + } + // scheduleAt must be in the future + scheduleAt := current.ScheduleAt.MustGet() + if !scheduleAt.After(time.Now().UTC()) { + return validate.WrapErrorWithField( + errors.New("schedule time must be in the future"), + "scheduleAt", + ) + } + // send_start_at must be more than 24h after schedule_at + if current.SendStartAt.IsSpecified() && !current.SendStartAt.IsNull() { + sendStartAt := current.SendStartAt.MustGet() + if sendStartAt.Sub(scheduleAt) < 24*time.Hour { + return validate.WrapErrorWithField( + errors.New("send start must be at least 24 hours after the schedule-at time"), + "scheduleAt", + ) + } + } + } // validate and update if err := current.Validate(); err != nil { return errs.Wrap(err) @@ -1638,13 +1727,27 @@ func (c *Campaign) UpdateByID( c.Logger.Errorw("failed to add recipient groups", "error", err) return errs.Wrap(err) } - // preserve jitter values from incoming campaign (not persisted to db) - current.JitterMin = incoming.JitterMin - current.JitterMax = incoming.JitterMax - err = c.schedule(ctx, session, current) - if err != nil { - c.Logger.Errorw("failed to re-schedule campaign", "error", err) - return errs.Wrap(err) + if current.ScheduleAt.IsSpecified() && !current.ScheduleAt.IsNull() { + // Late-scheduling: persist jitter to the DB so the task runner can apply it + // when schedule() is called hours later. Only overwrite jitter when the incoming + // payload explicitly specifies it — unspecified means "leave existing value alone". + if incoming.JitterMin.IsSpecified() { + current.JitterMin = incoming.JitterMin + current.JitterMax = incoming.JitterMax + } + err = c.setMostNotableCampaignEvent(ctx, current, data.EVENT_CAMPAIGN_PENDING_SCHEDULE) + if err != nil { + return errs.Wrap(err) + } + } else { + // Immediate scheduling: jitter lives in memory only for this call, never written to DB. + current.JitterMin = incoming.JitterMin + current.JitterMax = incoming.JitterMax + err = c.schedule(ctx, session, current) + if err != nil { + c.Logger.Errorw("failed to re-schedule campaign", "error", err) + return errs.Wrap(err) + } } c.AuditLogAuthorized(ae) return nil @@ -2906,6 +3009,65 @@ func (c *Campaign) HandleAnonymizeCampaigns( return nil } +// SchedulePendingCampaigns is called by the task runner. It finds all campaigns +// whose schedule_at time has passed and triggers schedule() for each. +func (c *Campaign) SchedulePendingCampaigns( + ctx context.Context, + session *model.Session, +) error { + ae := NewAuditEvent("Campaign.SchedulePendingCampaigns", session) + isAuthorized, err := IsAuthorized(session, data.PERMISSION_ALLOW_GLOBAL) + if err != nil && !errors.Is(err, errs.ErrAuthorizationFailed) { + c.LogAuthError(err) + return errs.Wrap(err) + } + if !isAuthorized { + c.AuditLogNotAuthorized(ae) + return errs.ErrAuthorizationFailed + } + pending, err := c.CampaignRepository.GetReadyToLateSchedule( + ctx, + &repository.CampaignOption{ + WithRecipientGroups: true, + WithAllowDeny: true, + }, + ) + if err != nil { + c.Logger.Errorw("failed to get campaigns ready for late scheduling", "error", err) + return errs.Wrap(err) + } + for _, campaign := range pending.Rows { + campaignID := campaign.ID.MustGet() + // Clear schedule_at BEFORE calling schedule() so that a concurrent task runner tick + // or a second server instance cannot pick up the same campaign and double-schedule it. + // If schedule() subsequently fails the campaign will remain in pending_schedule state + // with a null schedule_at; an operator can re-set schedule_at to retry. + clearScheduleAt := model.Campaign{} + clearScheduleAt.ScheduleAt.SetNull() + if err := c.CampaignRepository.UpdateByID(ctx, &campaignID, &clearScheduleAt); err != nil { + c.Logger.Errorw("failed to clear schedule_at before late-scheduling, skipping campaign", "campaignID", campaignID, "error", err) + // skip this campaign — better to retry next tick than to risk a double-schedule + continue + } + // Jitter is loaded from the DB columns (persisted at creation/update time). + if err := c.schedule(ctx, session, campaign); err != nil { + c.Logger.Errorw("failed to late-schedule campaign", "campaignID", campaignID, "error", err) + // continue to next — don't abort the whole run + continue + } + // Clear persisted jitter now that scheduling is done — it is no longer needed. + clearJitter := model.Campaign{} + clearJitter.JitterMin.SetNull() + clearJitter.JitterMax.SetNull() + if err := c.CampaignRepository.UpdateByID(ctx, &campaignID, &clearJitter); err != nil { + c.Logger.Errorw("failed to clear jitter after late-scheduling", "campaignID", campaignID, "error", err) + // non-fatal — jitter columns being non-null is harmless after scheduling + } + c.Logger.Infow("late-scheduled campaign", "campaignID", campaignID) + } + return nil +} + // CloseCampaignByID closes a campaign by id // DeleteDeviceCodesByCampaignID deletes all device codes for a campaign so every recipient // gets a fresh code (and picks up any proxy change) on their next page visit. diff --git a/backend/task/runner.go b/backend/task/runner.go index 35c0a24..1649355 100644 --- a/backend/task/runner.go +++ b/backend/task/runner.go @@ -195,6 +195,9 @@ func (d *Runner) ProcessSystemTasks( d.runTask("system - prune orphaned recipients", func() error { return d.PruneOrphanedRecipients(ctx, session) }) + d.runTask("system - late schedule campaigns", func() error { + return d.CampaignService.SchedulePendingCampaigns(ctx, session) + }) } // PruneOrphanedRecipients prunes orphaned recipients for global scope and all companies diff --git a/frontend/src/lib/api/api.js b/frontend/src/lib/api/api.js index 2e67418..4fa3379 100644 --- a/frontend/src/lib/api/api.js +++ b/frontend/src/lib/api/api.js @@ -513,6 +513,7 @@ export class API { * @param {string} campaign.sendEndAt * @param {string} [campaign.closeAt] * @param {string} [campaign.anonymizeAt] + * @param {string} [campaign.scheduleAt] * @param {string[]} campaign.recipientGroupIDs []uuid * @param {string[]} campaign.allowDenyIDs []uuid * @param {string} campaign.denyPageID uuid @@ -540,6 +541,7 @@ export class API { sendEndAt, closeAt, anonymizeAt, + scheduleAt, recipientGroupIDs, allowDenyIDs, denyPageID, @@ -566,6 +568,7 @@ export class API { sendEndAt, closeAt, anonymizeAt, + scheduleAt, recipientGroupIDs, allowDenyIDs, denyPageID, @@ -594,6 +597,7 @@ export class API { * @param {string} campaign.sendEndAt * @param {string} [campaign.closeAt] * @param {string} [campaign.anonymizeAt] + * @param {string} [campaign.scheduleAt] * @param {string} campaign.templateID uuid * @param {string[]} campaign.recipientGroupIDs []uuid * @param {string[]} campaign.allowDenyIDs []uuid @@ -622,6 +626,7 @@ export class API { sendEndAt, closeAt, anonymizeAt, + scheduleAt, recipientGroupIDs, allowDenyIDs, denyPageID, @@ -647,6 +652,7 @@ export class API { sendEndAt, closeAt, anonymizeAt, + scheduleAt, recipientGroupIDs, allowDenyIDs, denyPageID, diff --git a/frontend/src/lib/components/CheckboxField.svelte b/frontend/src/lib/components/CheckboxField.svelte index 5f28d40..a04cfcc 100644 --- a/frontend/src/lib/components/CheckboxField.svelte +++ b/frontend/src/lib/components/CheckboxField.svelte @@ -9,6 +9,7 @@ export let optional = false; export let id = null; export let inline = false; + export let disabled = false; let parentForm = null; let parentFormResetListener = null; @@ -44,6 +45,8 @@ class:flex-col={!inline} class:flex-row={inline} class:items-center={inline} + class:opacity-50={disabled} + class:cursor-not-allowed={disabled} >

@@ -65,7 +68,11 @@ {/if}

-