Re-resolve members absent after batch adds by forum user ID. Update renamed usernames and quarantine deleted users so they cannot block healthy group delivery. Align the fake with Discourse's partial-batch semantics.
405 lines
15 KiB
Go
405 lines
15 KiB
Go
// Package workflows holds the Discourse integration's Temporal workflows and
|
|
// activities: the periodic group-sync sweep (the level-triggered correctness
|
|
// backbone — design.md D3) and the targeted per-person reconcile that
|
|
// webhook events and link establishment trigger. Dispatch is direct-Temporal
|
|
// (FedWiki's model): the workflows own their database writes end-to-end, so
|
|
// the outbox handoff is not needed.
|
|
//
|
|
// Each activity is a thin transaction wrapper over a tx-bound logic
|
|
// function (convergeGroup, linkEntitledPersons, reconcilePerson) so the
|
|
// semantics are testable against a rolled-back transaction while production
|
|
// runs get per-activity commit boundaries.
|
|
package workflows
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
|
|
"git.coopcloud.tech/wiki-cafe/member-console/internal/integrations/discourse/client"
|
|
"git.coopcloud.tech/wiki-cafe/member-console/internal/integrations/discourse/linkage"
|
|
dcmod "git.coopcloud.tech/wiki-cafe/member-console/internal/integrations/discourse/store"
|
|
)
|
|
|
|
// groupClient is the slice of the API client the reconciler needs.
|
|
// AdminUser and VerifyKey serve the dead-link diagnosis path: usernames are
|
|
// a rename-unstable delivery address, so a member the batch add couldn't
|
|
// deliver is re-resolved by forum user id before the link is healed or
|
|
// quarantined.
|
|
type groupClient interface {
|
|
ListGroupMembers(ctx context.Context, groupName string) ([]client.User, error)
|
|
AddGroupMembers(ctx context.Context, groupID int64, usernames []string) error
|
|
RemoveGroupMembers(ctx context.Context, groupID int64, usernames []string) error
|
|
AdminUser(ctx context.Context, id int64) (*client.AdminUserRecord, error)
|
|
VerifyKey(ctx context.Context) error
|
|
}
|
|
|
|
// ActivitiesConfig wires the activities' dependencies. The adapter constructs
|
|
// it from viper reads (integration-owned config, not threaded through the
|
|
// generic WorkflowProvider parameters).
|
|
type ActivitiesConfig struct {
|
|
Database *sql.DB
|
|
Logger *slog.Logger
|
|
Client *client.Client
|
|
Linker *linkage.Linker
|
|
}
|
|
|
|
// Activities is the activity set registered against the shared worker.
|
|
type Activities struct {
|
|
cfg ActivitiesConfig
|
|
}
|
|
|
|
// NewActivities constructs the activity set.
|
|
func NewActivities(cfg ActivitiesConfig) *Activities {
|
|
return &Activities{cfg: cfg}
|
|
}
|
|
|
|
// inTx runs fn inside one transaction against the activity's DB.
|
|
func (a *Activities) inTx(ctx context.Context, fn func(q *dcmod.Queries) error) error {
|
|
tx, err := a.cfg.Database.BeginTx(ctx, nil)
|
|
if err != nil {
|
|
return fmt.Errorf("begin tx: %w", err)
|
|
}
|
|
defer tx.Rollback()
|
|
if err := fn(dcmod.New(tx)); err != nil {
|
|
return err
|
|
}
|
|
return tx.Commit()
|
|
}
|
|
|
|
// VerifyAPIKeyActivity fails the sweep up front when the admin API key is
|
|
// dead. On /admin/* routes an invalid key answers the same 404 as a missing
|
|
// record (live-verified, findings #9), so without this probe a revoked key
|
|
// degrades admin-record lookups into silent "not found"s instead of an
|
|
// error.
|
|
func (a *Activities) VerifyAPIKeyActivity(ctx context.Context) error {
|
|
return a.cfg.Client.VerifyKey(ctx)
|
|
}
|
|
|
|
// ListGroupMappingsActivity returns the operator-configured group mappings —
|
|
// the set of managed groups the sweep converges.
|
|
func (a *Activities) ListGroupMappingsActivity(ctx context.Context) ([]dcmod.GroupMapping, error) {
|
|
return dcmod.New(a.cfg.Database).ListGroupMappings(ctx)
|
|
}
|
|
|
|
// LinkStats reports a link pass's outcome.
|
|
type LinkStats struct {
|
|
Attempted int
|
|
Linked int
|
|
Unlinked int
|
|
Conflicts int
|
|
}
|
|
|
|
// LinkEntitledPersonsActivity attempts linkage for every entitled-but-
|
|
// unlinked person across the given resource keys.
|
|
func (a *Activities) LinkEntitledPersonsActivity(ctx context.Context, resourceKeys []string) (LinkStats, error) {
|
|
var stats LinkStats
|
|
err := a.inTx(ctx, func(q *dcmod.Queries) error {
|
|
var err error
|
|
stats, err = linkEntitledPersons(ctx, q, a.cfg.Linker, a.cfg.Logger, resourceKeys)
|
|
return err
|
|
})
|
|
return stats, err
|
|
}
|
|
|
|
// linkEntitledPersons attempts linkage for entitled-but-unlinked persons.
|
|
// Failures on individual persons are logged and skipped — an unreachable
|
|
// lookup for one person must not starve the rest of the pass; the next
|
|
// sweep retries.
|
|
func linkEntitledPersons(ctx context.Context, q *dcmod.Queries, linker *linkage.Linker, logger *slog.Logger, resourceKeys []string) (LinkStats, error) {
|
|
var stats LinkStats
|
|
seen := make(map[string]bool)
|
|
for _, key := range resourceKeys {
|
|
persons, err := q.ListEntitledUnlinkedPersons(ctx, key)
|
|
if err != nil {
|
|
return stats, fmt.Errorf("list entitled unlinked for %s: %w", key, err)
|
|
}
|
|
for _, p := range persons {
|
|
if seen[p.PersonID] {
|
|
continue
|
|
}
|
|
seen[p.PersonID] = true
|
|
stats.Attempted++
|
|
outcome, err := linker.EnsureLink(ctx, q, p.PersonID)
|
|
if err != nil {
|
|
logger.Warn("discourse linkage attempt failed",
|
|
slog.String("person_id", p.PersonID), slog.Any("error", err))
|
|
continue
|
|
}
|
|
switch outcome {
|
|
case linkage.OutcomeLinked, linkage.OutcomeCreatedLinked:
|
|
stats.Linked++
|
|
case linkage.OutcomeConflict:
|
|
stats.Conflicts++
|
|
default:
|
|
stats.Unlinked++
|
|
}
|
|
}
|
|
}
|
|
return stats, nil
|
|
}
|
|
|
|
// ConvergeStats reports one group's convergence outcome. Quarantined counts
|
|
// links orphaned this pass because their forum user no longer resolves.
|
|
type ConvergeStats struct {
|
|
MappingID string
|
|
GroupName string
|
|
Added int
|
|
Removed int
|
|
Desired int
|
|
Quarantined int
|
|
}
|
|
|
|
// ConvergeGroupActivity converges one managed group inside one transaction.
|
|
func (a *Activities) ConvergeGroupActivity(ctx context.Context, mapping dcmod.GroupMapping) (ConvergeStats, error) {
|
|
var stats ConvergeStats
|
|
err := a.inTx(ctx, func(q *dcmod.Queries) error {
|
|
var err error
|
|
stats, err = convergeGroup(ctx, q, a.cfg.Client, a.cfg.Logger, mapping)
|
|
return err
|
|
})
|
|
return stats, err
|
|
}
|
|
|
|
// convergeGroup diffs a managed group's observed membership against desired
|
|
// and converges Discourse: missing desired members are added, and any member
|
|
// not in the desired set is removed (managed groups are fully console-owned;
|
|
// forum-side drift is corrected and logged — spec: "Managed-group ownership
|
|
// boundary"). On success the observed projection is refreshed to the
|
|
// converged set.
|
|
func convergeGroup(ctx context.Context, q *dcmod.Queries, c groupClient, logger *slog.Logger, mapping dcmod.GroupMapping) (ConvergeStats, error) {
|
|
stats := ConvergeStats{MappingID: mapping.MappingID, GroupName: mapping.GroupName}
|
|
|
|
desired, err := q.ListDesiredLinkedMembers(ctx, mapping.ResourceKey)
|
|
if err != nil {
|
|
return stats, fmt.Errorf("list desired members: %w", err)
|
|
}
|
|
stats.Desired = len(desired)
|
|
desiredByID := make(map[int64]string, len(desired))
|
|
for _, d := range desired {
|
|
desiredByID[d.DiscourseUserID] = d.DiscourseUsername
|
|
}
|
|
|
|
observed, err := c.ListGroupMembers(ctx, mapping.GroupName)
|
|
if err != nil {
|
|
return stats, fmt.Errorf("list group members: %w", err)
|
|
}
|
|
observedByID := make(map[int64]string, len(observed))
|
|
for _, m := range observed {
|
|
observedByID[m.ID] = m.Username
|
|
}
|
|
|
|
missing := make(map[int64]dcmod.ListDesiredLinkedMembersRow)
|
|
for _, d := range desired {
|
|
if _, ok := observedByID[d.DiscourseUserID]; !ok {
|
|
missing[d.DiscourseUserID] = d
|
|
}
|
|
}
|
|
var removeNames []string
|
|
for id, username := range observedByID {
|
|
if _, ok := desiredByID[id]; !ok {
|
|
removeNames = append(removeNames, username)
|
|
}
|
|
}
|
|
|
|
undelivered := make(map[int64]bool)
|
|
if len(missing) > 0 {
|
|
var err error
|
|
undelivered, err = deliverMissingMembers(ctx, q, c, logger, mapping, missing, &stats)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
if stats.Added > 0 {
|
|
logger.Info("discourse group membership converged",
|
|
slog.String("group", mapping.GroupName), slog.Int("added", stats.Added))
|
|
}
|
|
}
|
|
if len(removeNames) > 0 {
|
|
if err := c.RemoveGroupMembers(ctx, mapping.DiscourseGroupID, removeNames); err != nil {
|
|
return stats, fmt.Errorf("remove members from %s: %w", mapping.GroupName, err)
|
|
}
|
|
stats.Removed = len(removeNames)
|
|
// Drift correction is logged with the discrepancy, per spec.
|
|
logger.Info("discourse managed-group drift corrected",
|
|
slog.String("group", mapping.GroupName), slog.Any("removed", removeNames))
|
|
}
|
|
|
|
// Refresh the observed projection to the converged set — desired minus
|
|
// the members this pass could not deliver (quarantined or deferred to
|
|
// the next sweep): the projection records reality, not intent.
|
|
if err := q.ReplaceObservedGroupMembersDelete(ctx, mapping.MappingID); err != nil {
|
|
return stats, fmt.Errorf("clear observed projection: %w", err)
|
|
}
|
|
for id := range desiredByID {
|
|
if undelivered[id] {
|
|
continue
|
|
}
|
|
if err := q.InsertObservedGroupMember(ctx, dcmod.InsertObservedGroupMemberParams{
|
|
MappingID: mapping.MappingID,
|
|
DiscourseUserID: id,
|
|
}); err != nil {
|
|
return stats, fmt.Errorf("insert observed member: %w", err)
|
|
}
|
|
}
|
|
return stats, nil
|
|
}
|
|
|
|
// deliverMissingMembers adds the desired-but-absent members. The batched
|
|
// add is the common path, but usernames are a rename-unstable delivery
|
|
// address and a link can outlive its forum account: live 3.5.3 silently
|
|
// drops unresolvable names from a partly-valid batch and rejects an
|
|
// all-unresolvable batch with a 400 — either way a dead link would go
|
|
// undelivered forever or, in the all-dead case, fail every subsequent sweep
|
|
// (the 2026-07-22 "one dead username poisons the whole batch add"
|
|
// incident). Members still absent after the batch attempt are therefore
|
|
// diagnosed individually against the forum's id-keyed admin record:
|
|
// renamed users heal the stored username and retry; users the forum no
|
|
// longer knows are quarantined as 'orphaned', leaving the desired set for
|
|
// good. Transient per-member failures are logged and left for the next
|
|
// sweep — one sick member must not starve the rest (the
|
|
// linkEntitledPersons isolation principle). Returns the ids that remain
|
|
// undelivered so the observed projection records only reality.
|
|
func deliverMissingMembers(ctx context.Context, q *dcmod.Queries, c groupClient, logger *slog.Logger, mapping dcmod.GroupMapping, missing map[int64]dcmod.ListDesiredLinkedMembersRow, stats *ConvergeStats) (map[int64]bool, error) {
|
|
names := make([]string, 0, len(missing))
|
|
for _, m := range missing {
|
|
names = append(names, m.DiscourseUsername)
|
|
}
|
|
if err := c.AddGroupMembers(ctx, mapping.DiscourseGroupID, names); err != nil {
|
|
if !client.IsUnresolvableUsernames(err) {
|
|
return nil, fmt.Errorf("add members to %s: %w", mapping.GroupName, err)
|
|
}
|
|
logger.Warn("discourse batch add unresolvable; diagnosing members individually",
|
|
slog.String("group", mapping.GroupName), slog.Int("batch", len(names)))
|
|
}
|
|
|
|
// The batch response reports nothing per-name, so re-list to see what
|
|
// actually landed before diagnosing the remainder.
|
|
current, err := c.ListGroupMembers(ctx, mapping.GroupName)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("re-list group members: %w", err)
|
|
}
|
|
inGroup := make(map[int64]bool, len(current))
|
|
for _, m := range current {
|
|
inGroup[m.ID] = true
|
|
}
|
|
|
|
undelivered := make(map[int64]bool)
|
|
for id, row := range missing {
|
|
if inGroup[id] {
|
|
stats.Added++
|
|
continue
|
|
}
|
|
rec, err := c.AdminUser(ctx, id)
|
|
if err != nil {
|
|
logger.Warn("discourse member diagnosis failed; next sweep retries",
|
|
slog.Int64("discourse_user_id", id), slog.Any("error", err))
|
|
undelivered[id] = true
|
|
continue
|
|
}
|
|
if rec == nil {
|
|
// 404 — but a dead admin key answers the same on /admin/*
|
|
// (findings #9), so re-verify the key before condemning the
|
|
// link. A key failing here fails the converge: every remaining
|
|
// diagnosis would be equally blind.
|
|
if err := c.VerifyKey(ctx); err != nil {
|
|
return nil, fmt.Errorf("verify key during orphan check: %w", err)
|
|
}
|
|
if _, err := q.SetUserLinkOrphaned(ctx, row.PersonID); err != nil {
|
|
return nil, fmt.Errorf("quarantine link for person %s: %w", row.PersonID, err)
|
|
}
|
|
stats.Quarantined++
|
|
undelivered[id] = true
|
|
logger.Warn("discourse link quarantined: forum user no longer exists",
|
|
slog.String("person_id", row.PersonID),
|
|
slog.Int64("discourse_user_id", id),
|
|
slog.String("username", row.DiscourseUsername),
|
|
slog.String("group", mapping.GroupName))
|
|
continue
|
|
}
|
|
if rec.Username != row.DiscourseUsername {
|
|
if _, err := q.UpdateUserLinkUsername(ctx, dcmod.UpdateUserLinkUsernameParams{
|
|
PersonID: row.PersonID,
|
|
DiscourseUsername: rec.Username,
|
|
}); err != nil {
|
|
return nil, fmt.Errorf("heal username for person %s: %w", row.PersonID, err)
|
|
}
|
|
logger.Info("discourse link username healed",
|
|
slog.String("person_id", row.PersonID),
|
|
slog.String("old", row.DiscourseUsername),
|
|
slog.String("new", rec.Username))
|
|
}
|
|
if err := c.AddGroupMembers(ctx, mapping.DiscourseGroupID, []string{rec.Username}); err != nil {
|
|
logger.Warn("discourse individual add failed; next sweep retries",
|
|
slog.String("username", rec.Username), slog.Any("error", err))
|
|
undelivered[id] = true
|
|
continue
|
|
}
|
|
stats.Added++
|
|
}
|
|
return undelivered, nil
|
|
}
|
|
|
|
// ReconcilePersonActivity runs the targeted fast path in one transaction.
|
|
func (a *Activities) ReconcilePersonActivity(ctx context.Context, personID string) error {
|
|
return a.inTx(ctx, func(q *dcmod.Queries) error {
|
|
return reconcilePerson(ctx, q, a.cfg.Client, a.cfg.Linker, personID)
|
|
})
|
|
}
|
|
|
|
// reconcilePerson ensures the person's link, then converges exactly their
|
|
// membership in every managed group.
|
|
func reconcilePerson(ctx context.Context, q *dcmod.Queries, c groupClient, linker *linkage.Linker, personID string) error {
|
|
if _, err := linker.EnsureLink(ctx, q, personID); err != nil {
|
|
return fmt.Errorf("ensure link: %w", err)
|
|
}
|
|
link, err := q.GetUserLinkByPersonID(ctx, personID)
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
// Unlinked — nothing to converge; the sweep remains the safety net.
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return fmt.Errorf("get link: %w", err)
|
|
}
|
|
if link.Status != "linked" {
|
|
return nil
|
|
}
|
|
|
|
mappings, err := q.ListGroupMappings(ctx)
|
|
if err != nil {
|
|
return fmt.Errorf("list group mappings: %w", err)
|
|
}
|
|
for _, mapping := range mappings {
|
|
desired, err := q.IsPersonDesiredForResourceKey(ctx, dcmod.IsPersonDesiredForResourceKeyParams{
|
|
ResourceKey: mapping.ResourceKey,
|
|
PersonID: personID,
|
|
})
|
|
if err != nil {
|
|
return fmt.Errorf("desired check for %s: %w", mapping.GroupName, err)
|
|
}
|
|
if desired {
|
|
if err := c.AddGroupMembers(ctx, mapping.DiscourseGroupID, []string{link.DiscourseUsername}); err != nil {
|
|
return fmt.Errorf("add %s to %s: %w", link.DiscourseUsername, mapping.GroupName, err)
|
|
}
|
|
if err := q.InsertObservedGroupMember(ctx, dcmod.InsertObservedGroupMemberParams{
|
|
MappingID: mapping.MappingID,
|
|
DiscourseUserID: link.DiscourseUserID,
|
|
}); err != nil {
|
|
return fmt.Errorf("record observed member: %w", err)
|
|
}
|
|
} else {
|
|
if err := c.RemoveGroupMembers(ctx, mapping.DiscourseGroupID, []string{link.DiscourseUsername}); err != nil {
|
|
return fmt.Errorf("remove %s from %s: %w", link.DiscourseUsername, mapping.GroupName, err)
|
|
}
|
|
if _, err := q.RemoveObservedGroupMember(ctx, dcmod.RemoveObservedGroupMemberParams{
|
|
MappingID: mapping.MappingID,
|
|
DiscourseUserID: link.DiscourseUserID,
|
|
}); err != nil {
|
|
return fmt.Errorf("clear observed member: %w", err)
|
|
}
|
|
}
|
|
}
|
|
return nil
|
|
}
|