-
Notifications
You must be signed in to change notification settings - Fork 34
Expand file tree
/
Copy pathservice_worker_fence_postgres_test.go
More file actions
83 lines (76 loc) · 2.67 KB
/
Copy pathservice_worker_fence_postgres_test.go
File metadata and controls
83 lines (76 loc) · 2.67 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
package kata
import (
"context"
"database/sql"
"errors"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.kenn.io/kata/internal/db"
"go.kenn.io/kata/internal/testenv"
)
func TestServicePostgresWorkerTransactionFenceRollsBackAndStopsWorkers(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancel()
dsn, cleanup := testenv.NewPostgresContainer(t, ctx)
t.Cleanup(cleanup)
rejected := errors.New("worker transaction rejected")
service, err := New(ctx, Config{
DSN: dsn,
Postgres: PostgresConfig{Schema: "kata", SchemaMode: PostgresSchemaBootstrap},
Access: allowWorkerFenceAccess{},
WorkerTransactionFence: func(ctx context.Context, tx Transaction) error {
if _, err := tx.ExecContext(ctx,
`INSERT INTO kata.worker_fence_markers(attempt) VALUES($1)`, 1); err != nil {
return err
}
return rejected
},
})
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, service.Close()) })
project, err := service.store.CreateProject(ctx, "worker-project")
require.NoError(t, err)
_, err = service.store.EnableProjectFederation(ctx, project.ID, "worker")
require.NoError(t, err)
issue, _, err := service.store.CreateIssue(ctx, db.CreateIssueParams{
ProjectID: project.ID, Title: "expired claim", Author: "worker",
})
require.NoError(t, err)
_, err = service.store.AcquireClaim(ctx, db.AcquireClaimParams{
ProjectID: project.ID,
IssueRef: issue.ShortID,
Principal: db.ClaimPrincipal{
HolderInstanceUID: service.store.InstanceUID(),
Holder: "worker",
ClientKind: "test",
},
ClaimKind: "timed",
TTL: time.Minute,
Now: time.Now().UTC().Add(-2 * time.Minute),
})
require.NoError(t, err)
inspection, err := sql.Open("pgx", dsn)
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, inspection.Close()) })
_, err = inspection.ExecContext(ctx,
`CREATE TABLE kata.worker_fence_markers (attempt BIGINT NOT NULL)`)
require.NoError(t, err)
runDone := make(chan error, 1)
go func() { runDone <- service.Run(ctx) }()
select {
case runErr := <-runDone:
require.ErrorIs(t, runErr, rejected)
case <-time.After(10 * time.Second):
require.FailNow(t, "worker fence rejection did not stop the service workers")
}
var markerCount, released int
require.NoError(t, inspection.QueryRowContext(ctx,
`SELECT count(*) FROM kata.worker_fence_markers`).Scan(&markerCount))
require.NoError(t, inspection.QueryRowContext(ctx,
`SELECT CASE WHEN released_at IS NOT NULL THEN 1 ELSE 0 END
FROM kata.issue_claims WHERE issue_uid = $1`, issue.UID).Scan(&released))
assert.Zero(t, markerCount)
assert.Zero(t, released)
}