After DarkKranti
FireplaceA productivity app pulled out of a Go monolith, one gRPC service at a time.

Published to nobody

I built a transactional outbox so plan-service could never lose an event and an inbox so insights-service could never process one twice, then never plugged them together.

New moon, 1% lit · new moonBy Kranti · 10 min readgo · postgres · rabbitmq · outbox · distributed-systems

Fireplace has a transactional outbox, and it's some of the most careful plumbing in the repo. A plan and its plan.created event commit in one Postgres transaction. A worker drains pending events to RabbitMQ with FOR UPDATE SKIP LOCKED. On the far side, a consumer dedups on a primary key so a redelivered event is a no-op. I wrote the producer half in June, the consumer half the week after, and felt pretty good about it.

The worker stamps published_at on every event it drains. None of those events has ever reached a consumer. Both of those sentences are true, and the outbox isn't the broken part.

Ten databases and no shared transaction

Setup first. Fireplace started life as one Go/Gin monolith and is being pulled apart with the strangler pattern. The gateway stays the public HTTP surface, and domains move out into gRPC services one slice at a time. Each service with data of its own got its own Postgres. At the peak, compose ran ten Postgres containers, a production and a test database per service, which is noise on a laptop and unaffordable on the 1-2GB VPS this is meant to ship on. ADR-0010 collapsed them into one instance with three databases. The one hard rule it kept is that plan-service and insights-service never share a database.

That rule is why this post exists. If plans and insights lived in one database, "a plan was created, go generate suggestions for it" could be a single transaction, or a join. They don't, on purpose. The only way that fact crosses the boundary is a message. And once a service has to write to its database and also send a message, it's updating two systems with no transaction that spans both.

Insert, then publish, then hope

The naive version still lives in the repo, one package over. When you tick a checklist item, plan-service updates the row, reloads it, and calls this:

services/plan-service/internal/checklistitem/publisher.gogo
func (s *service) PublishItemCompleted(ctx context.Context, item *Item) {
	body, err := proto.Marshal(&eventspb.ChecklistItemCompletedEvent{
		Id:          item.ID.String(),
		PlanId:      item.PlanID.String(),
		CompletedAt: timestamppb.Now(),
	})
	if err != nil {
		slog.ErrorContext(ctx, "failed to marshal checklist_item.completed", "err", err, "item_id", item.ID)
		return
	}
	if err := s.publishCh.PublishWithContext(ctx,
		commonconstants.PlanEventsExchange,
		commonconstants.ChecklistItemCompleted,
		commonbroker.Message{
			ContentType:  "application/protobuf",
			Body:         body,
			DeliveryMode: commonbroker.Persistent,
		},
	); err != nil {
		slog.ErrorContext(ctx, "failed to publish checklist_item.completed", "err", err, "item_id", item.ID)
	}
}

Look at the error branch. If the publish fails, the item is done in the database, the event is gone, and the only record is a log line. There's no retry, and nothing to retry from.

That's the dual-write problem. It comes in two shapes:

  1. Commit, then publish. The process dies or the broker blips between the two. The row exists and the event never will. Downstream quietly misses it.
  2. Publish, then commit. The event goes out, then the commit fails. Now a consumer is acting on a row that was rolled back.

Ordering can't save you. Each order just picks how you lose. Nothing subscribes to checklist_item.completed yet, so today this costs nothing. plan.created was supposed to feed insights-service, so plans got the outbox.

The outbox

The fix is to stop doing two writes. Do one write, to your own database, and let a separate process turn it into a message later. The event becomes a row:

services/plan-service/migrations/000020_create_outbox.up.sqlsql
CREATE TABLE outbox (
    -- also used as event_id
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
 
    -- Core routing and event metadata
    routing_key VARCHAR(255) NOT NULL,
    exchange VARCHAR(255) NOT NULL,
    payload BYTEA NOT NULL,
 
    -- Timestamps
    created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
 
    -- State of published
    -- null = pending
    -- not null (existing timestamp) = published
    -- whole row not present = no issues
    published_at TIMESTAMPTZ NULL
);
 
CREATE INDEX idx_outbox_pending ON outbox(created_at) WHERE published_at IS NULL;

A null published_at means pending. The index is partial, so it only covers pending rows and stays small however much history piles up, and the drain can walk it in created_at order.

Create then writes the plan and the event row in the same transaction:

services/plan-service/internal/plan/service.gogo
func (s *service) Create(ctx context.Context, in *CreatePlanInput) (*Plan, error) {
	p := &Plan{UserID: in.UserID, Name: in.Name, Focus: in.Focus, PlanType: in.PlanType}
	var newPlan *Plan
	err := commonhelpers.ExecTx(ctx, s.db, func(tx *sqlx.Tx) error {
		plan, err := s.repo.CreateTx(ctx, tx, p)
		if err != nil {
			return err
		}
		payload, err := proto.Marshal(&pbevents.PlanCreatedEvent{
			Id:        plan.ID.String(),
			UserId:    plan.UserID.String(),
			Name:      plan.Name,
			Focus:     plan.Focus,
			PlanType:  plan.PlanType,
			CreatedAt: timestamppb.New(plan.CreatedAt),
		})
		if err != nil {
			return fmt.Errorf("create plan error when marshalling to proto plan id %s: %w", plan.ID, err)
		}
		err = s.outboxService.CreateTx(ctx, tx, outbox.CreateOutboxParams{
			RoutingKey: commonconstants.PlanCreated,
			Exchange:   commonconstants.PlanEventsExchange,
			Payload:    payload,
		})
		if err != nil {
			return err
		}
		newPlan = plan
		return nil
	})
	if err != nil {
		return nil, err
	}
	return newPlan, nil
}

ExecTx in common/utils begins a transaction, runs the closure, rolls back on an error or a panic, and commits otherwise. So the plan and its event both exist or neither does. If the proto marshal fails, the plan rolls back with it. Nothing is sent here. The event is a promise sitting in a table, and the promise commits atomically with the thing it describes.

The old direct publisher, PublishPlanCreated, is still in plan/publisher.go with its publish call commented out and nothing calling it. I keep meaning to delete it.

Draining it

The worker lives in common/worker, because the plan is for every producer to get one. It takes an EventDrainer (anything that can list pending events and mark them published), a publisher, and a poll interval. plan-service builds it with time.Minute*2. Each tick runs one drain:

common/worker/publishworker.gogo
func (w *PublishWorker) Drain(ctx context.Context) error {
	return commonhelpers.ExecTx(ctx, w.db, func(tx *sqlx.Tx) error {
		events, err := w.eventDrainer.GetUnpublished(ctx, tx)
		if err != nil {
			return fmt.Errorf("worker draining unpublished events: %w", err)
		}
 
		successfulIds := make([]uuid.UUID, 0, len(events))
		for _, event := range events {
			err := w.publisher.PublishWithContext(ctx, event.Exchange, event.RoutingKey,
				commonbroker.Message{
					MessageId:    event.ID.String(),
					ContentType:  "application/protobuf",
					Body:         event.Payload,
					DeliveryMode: commonbroker.Persistent,
				})
			if err != nil {
				// couldn't publish, leave for next worker
				slog.Warn("worker publish attempt failed for event", "event_id", event.ID, "error", err)
				continue
			}
			successfulIds = append(successfulIds, event.ID)
		}
		if len(successfulIds) == 0 {
			return nil
		}
 
		err = w.eventDrainer.MarkPublished(ctx, tx, successfulIds)
		if err != nil {
			return fmt.Errorf("worker couldn't mark published: %w", err)
		}
		return nil
	})
}

GetUnpublished is where the interesting SQL lives:

services/plan-service/internal/outbox/repository.gosql
SELECT
    id,
    routing_key,
    exchange,
    payload,
    published_at,
    created_at
FROM outbox
WHERE published_at IS NULL
ORDER BY created_at ASC
LIMIT $1
FOR UPDATE SKIP LOCKED

$1 is defaultUnpublishedLimit, which is 10. FOR UPDATE locks the rows it read until the transaction ends. SKIP LOCKED means a second worker running the same query doesn't wait on those rows. It skips them and takes the next ten. So two plan-service instances can drain one table without publishing the same row twice. Today there's one instance on one box, so it's insurance, but cheap insurance.

The numbers deserve saying out loud. Ten rows every two minutes is a ceiling of 5 events a minute, 300 an hour. A time.Ticker doesn't fire on start, so the first drain after a deploy happens two minutes in, and a fresh event can sit for up to two minutes before anyone looks. On shutdown, Run does one last drain on a new context with a 5-second timeout, so a SIGTERM doesn't strand whatever committed since the last tick. For Fireplace's rate of plan creation this is all fine. It's also a number I picked once and never looked at again.

Which is why the MessageId line matters. Every message carries the outbox row's UUID. A republished event has the same id as the first copy, and the consumer's whole job is to notice.

The other end

insights-service has a processed_events table with PRIMARY KEY (event_id, consumer). Handling an event means inserting that pair and writing the insight in one transaction. A second delivery hits the primary key, which is supposed to surface as ErrDuplicateResource, and the consumer acks and drops it:

services/insights-service/internal/insights/service.gogo
func (s *Service) Create(ctx context.Context, param CreateInsightFromPlanParam) error {
	key := fmt.Sprintf("dedup:insights:%s", param.EventID)
	acquired, _ := s.cache.SetNX(ctx, key, inProgressMarker, time.Second*5).Result()
	if !acquired {
		return ErrEventAlreadyProcessed
	}
 
	if err := commonhelpers.ExecTx(ctx, s.db, func(tx *sqlx.Tx) error {
		// attempt to create dedup table first, rollback on conflict, true authority here
		if err := s.inboxService.CreateTx(ctx, tx, param.EventID); err != nil {
			if errors.Is(err, commonconstants.ErrDuplicateResource) {
				return fmt.Errorf("event %s: %w", param.EventID, ErrEventAlreadyProcessed)
			}
			return err
		}
		return s.repo.CreateTx(ctx, tx, CreateInsightParam{
			PlanID:      param.PlanID,
			UserID:      param.UserID,
			InsightType: string(InsightTypeSuggestion),
			Content:     "", // TODO: need to acquire and fill
		})
	}); err != nil {
		s.cache.Del(ctx, key)
		return err
	}
 
	s.cache.Set(ctx, key, inboxWriteComplete, time.Hour*24)
	return nil
}

The Redis SetNX in front is a 5-second claim while the event is in flight, bumped to 24 hours once the transaction commits. It's meant as an efficiency layer so two racing deliveries don't both reach Postgres. The primary key is the authority. ADR-0008 calls this effectively-once on every hop: outbox on the producer, inbox on the consumer, and the inbox row commits with the business write or not at all. (Yes, Content is an empty string. That's its own post.)

On paper that closes the loop. At-least-once delivery plus an idempotent consumer gives you exactly-once effects.

Nobody home

This is SetupServices for insights-service, the composition root where everything gets wired:

services/insights-service/config/services.gogo
func SetupServices(db *sqlx.DB, _ *amqp.Channel, registry discovery.Registry, redisClient *redis.Client) (*grpc.Server, error) {
	planClient := insights.NewPlanClient(registry)
 
	checklistGen := ai.NewChecklistGen()
	searchTermGen := ai.NewSearchTermGen()
 
	youtubeFinder, err := videodiscovery.NewYoutubeVideoFinder()
	if err != nil {
		return nil, fmt.Errorf("insights: init youtube video finder: %w", err)
	}
	videoFinder := insights.NewDiscoveryVideoFinder(youtubeFinder)
 
	repo := insights.NewRepository(db)
	service := insights.NewService(planClient, checklistGen, searchTermGen, videoFinder, redisClient, repo, db)
	handler := insights.NewHandler(service)
 
	grpcServer := grpc.NewServer()
	pb.RegisterInsightsServiceServer(grpcServer, handler)
 
	slog.Info("insights-service initialized successfully")
	return grpcServer, nil
}

Second parameter. The AMQP channel comes in and gets thrown away with a _. The consumer is never built, so it never subscribes. It couldn't be started anyway. There's no Listen method, and consumePlanEvents is unexported and called by nothing. I'd been describing this as "the pipeline ends in a queue nobody reads". Checking the code for this post, it's worse. There is no queue.

The only code that binds anything to plan.created is the consumer's SetupAMQPInfrastructure, and nothing calls it. If something did, it would declare plan-service.events (plan-service's own queue, copy-pasted along with log lines that still say plan-service:) and then bind insights-service.plan.created, which nothing ever declares, so the bind would fail. Meanwhile the publisher calls ch.Publish(exchange, key, false, false, ...). Those two falses are mandatory and immediate. With mandatory off and no queue bound to the routing key, RabbitMQ accepts the message and drops it at the exchange. Publish returns nil, the id goes into successfulIds, and MarkPublished stamps the row.

So the outbox is a tidy ledger of events that were atomically recorded, reliably drained and delivered to no one. Every guarantee held, and every event was dropped anyway.

It wouldn't have worked on the day it started, either. Ticket I-0026 lists what stands between this and a running consumer. inboxService is never passed to NewService, so it's a nil interface and the first real event panics on s.inboxService.CreateTx, right in the exactly-once path. The Redis claim discards its error, so with Redis down acquired is false, which reads as "already processed", which the consumer acks. The comment says best-effort; the code fails closed. And FS-0006 has since moved the trigger to plan.items_requested, leaving plan.created as a plain lifecycle fact.

Rereading it for this post I found two more the ticket doesn't have. On the success path the consumer never acks. Create returns nil, the case ends, and the delivery stays unacked. And the duplicate check in the snippet above can't fire. WrapDBErr spots a unique violation by matching pgx's *pgconn.PgError, but the services connect through lib/pq, whose errors are *pq.Error. A real duplicate would come back as ErrUnexpectedError, get nacked without requeue, and vanish.1 Right outcome, wrong reason, wrong log line.

What I actually got wrong wasn't the outbox. I built both halves, tested neither end to end, and a producer with no consumer looks exactly like a healthy one from where the producer sits. published_at gets set either way. I-0026 is marked human-owned, which means it's mine, and it's still open.

Notes

  1. The poison paths all say "DLQ" in their comments and call Nack(false, false), but no queue in the service is declared with a dead-letter exchange. A rejected message is just discarded. ↩