1 Commits
Author SHA1 Message Date
clawbot fa67d883e9 Close archive writers when the delivery engine stops (closes #280)
check / check (push) Successful in 3m36s
The engine cached archive writers and never closed them at shutdown,
so after a clean stop an archive's rows could sit in its -wal while
the .db held no table. The engine's stop hook now evicts every cached
writer once its workers have returned, the same way deleting a webhook
does, so a clean stop leaves each archive as one file and a late write
is refused. If the workers do not return within the stop budget, the
writers are left open as a kill would leave them.

The README no longer says archives keep their sidecars across a clean
stop.

Model: opus-5-5
2026-09-29 05:00:30 +00:00
6 changed files with 18 additions and 113 deletions
+3 -3
View File
@@ -368,9 +368,9 @@ func (e *Engine) start() {
// writer for long by then: the archive sweeper stops before the
// engine, and deleting a webhook only closes one. If the pool did
// not drain in time, the writers are left open, as a kill would
// leave them. Closing them would wait for any write in progress,
// and a worker still running would then open new writers that
// nothing closes, so it gains nothing over a kill.
// leave them: a worker still running may be mid-write, and
// closing its writer would wait on that write and then fail the
// next delivery the worker archives.
func (e *Engine) stop(ctx context.Context) error {
e.log.Info("delivery engine stopping")
+4 -4
View File
@@ -327,10 +327,10 @@ func TestEngine_StopHookClosesArchives(t *testing.T) {
}
// TestEngine_StopHookTimeoutLeavesArchivesOpen covers a stop whose
// budget runs out while a worker is still running. The archive
// writers are left open, as a kill would leave them: closing them
// would wait for any write in progress, and that worker would then
// open new writers that nothing closes.
// budget runs out while a worker is still running. That worker may
// be in the middle of an archive write, so the archive writers are
// left open, as a kill would leave them, rather than closed
// underneath it.
func TestEngine_StopHookTimeoutLeavesArchivesOpen(t *testing.T) {
t.Parallel()
-79
View File
@@ -1166,10 +1166,6 @@ func TestIsForwardableHeader(t *testing.T) {
assert.False(t,
delivery.ExportIsForwardableHeader("Content-Length"),
)
assert.False(t,
delivery.ExportIsForwardableHeader("Content-Type"),
)
}
func TestTruncate(t *testing.T) {
@@ -1251,81 +1247,6 @@ func TestDoHTTPRequest_ForwardsHeaders(t *testing.T) {
)
}
// The event's stored inbound headers carry the same Content-Type the
// receiver saved as the event's ContentType, so a delivery could send
// it twice. It must go out exactly once, with a Content-Type configured
// on the target winning, then the event's ContentType.
func TestApplyRequestHeaders_SendsOneContentType(t *testing.T) {
t.Parallel()
cases := map[string]struct {
inbound string
event string
configured string
want []string
}{
"inbound and event agree": {
inbound: testContentType,
event: testContentType,
want: []string{testContentType},
},
"inbound and event disagree": {
inbound: "text/plain",
event: testContentType,
want: []string{testContentType},
},
"event has none": {
inbound: testContentType,
want: nil,
},
"target configures its own": {
inbound: testContentType,
event: testContentType,
configured: "application/xml",
want: []string{"application/xml"},
},
}
for name, tc := range cases {
t.Run(name, func(t *testing.T) {
t.Parallel()
inbound, err := json.Marshal(map[string][]string{
headerContentType: {tc.inbound},
})
require.NoError(t, err)
cfg := &delivery.HTTPTargetConfig{}
if tc.configured != "" {
cfg.Headers = map[string]string{
headerContentType: tc.configured,
}
}
req, err := http.NewRequestWithContext(
context.Background(),
http.MethodPost,
"https://target.example.com/hook",
http.NoBody,
)
require.NoError(t, err)
delivery.ExportApplyRequestHeaders(
req,
&database.Event{
Headers: string(inbound),
ContentType: tc.event,
},
cfg,
)
assert.Equal(t,
tc.want, req.Header.Values(headerContentType),
)
})
}
}
func TestProcessDelivery_RoutesToCorrectHandler(
t *testing.T,
) {
+6 -12
View File
@@ -339,11 +339,10 @@ func TestRedirectPolicy_StopsAtHopCap(t *testing.T) {
// The set the redirect policy strips is whatever the delivery path
// actually put on the wire, so a header added to the forward set is
// covered without a second edit. A header the event never carried
// is not in the set, and neither is the inbound Content-Type, because
// it is not forwarded. Two more are deliberately excluded: a
// Content-Type configured on the target describes the body, which a
// 307 carries across hosts, and the inbound User-Agent every real
// sender supplies is overwritten before the request goes out.
// is not in the set, and the delivery path's own two are deliberately
// excluded: Content-Type describes the body, which a 307 carries
// across hosts, and the inbound User-Agent every real sender supplies
// is overwritten before the request goes out.
func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
t.Parallel()
@@ -372,7 +371,6 @@ func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
&delivery.HTTPTargetConfig{
Headers: map[string]string{
probeHeaderName: probeHeaderValue,
"Content-Type": testContentType,
},
},
)
@@ -380,11 +378,7 @@ func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
assert.Equal(t,
[]string{probeHeaderName, inboundHeaderName}, names,
"both header classes are reported, and only those: "+
"Host and the inbound Content-Type are never "+
"forwarded, User-Agent is the delivery path's own",
)
assert.NotContains(t, names, "Content-Type",
"a Content-Type configured on the target must survive "+
"a cross-origin 307/308 with the body it describes",
"Host is never forwarded, Content-Type and "+
"User-Agent are the delivery path's own",
)
}
+1 -2
View File
@@ -11,11 +11,10 @@ import (
"sneak.berlin/go/webhooker/internal/delivery"
)
// Literals these tests repeat, named so that the header names and the
// Literals these tests repeat, named so that the header name and the
// keep-forever archive config each have one definition.
const (
headerAuthorization = "Authorization"
headerContentType = "Content-Type"
bearerValue = "Bearer abc"
archiveConfigNever = "{\"expiry\":\"never\"}"
)
+4 -13
View File
@@ -541,11 +541,6 @@ func isForwardableHeader(name string) bool {
"Upgrade", "Proxy-Authorization",
"Proxy-Connection", "Content-Length":
return false
case "Content-Type":
// applyRequestHeaders sets Content-Type itself. The receiver
// already stored this inbound value as the event's
// ContentType, so forwarding it too would send it twice.
return false
default:
return true
}
@@ -558,10 +553,6 @@ func isForwardableHeader(name string) bool {
// policy strips exactly that set on a hop that leaves the origin,
// so the forward set is decided here and only here — a header added
// to it is covered off-origin without a second edit elsewhere.
//
// Content-Type goes out once: a Content-Type configured on the target
// wins, otherwise the event's ContentType, otherwise none. The inbound
// Content-Type in the event's headers is never forwarded.
func applyRequestHeaders(
req *http.Request,
event *database.Event,
@@ -582,10 +573,10 @@ func applyRequestHeaders(
req.Header.Set("User-Agent", "webhooker/1.0")
// A Content-Type configured on the target describes the body
// being sent rather than the sender. A 307/308 preserves the
// body across hosts, so stripping it would send that body
// untyped.
// Content-Type describes the body being sent rather than the
// sender, and the delivery path sets it from the event itself.
// A 307/308 preserves the body across hosts, so stripping it
// would send that body untyped.
delete(originScoped, "Content-Type")
// User-Agent is overwritten just above, so an inbound one never