Files
T
cgalo5758 5fb9ba8dfe Recover from stale Discourse user links
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.
2026-07-22 14:02:40 -05:00

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
}