-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy patha.patch
More file actions
1956 lines (1869 loc) · 73.7 KB
/
Copy patha.patch
File metadata and controls
1956 lines (1869 loc) · 73.7 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
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
commit 1e47627c350451039b1fc801186365320a50409f
Author: Paweł Gronowski <pawel.gronowski@docker.com>
Date: Wed Aug 12 19:05:19 2026 +0200
extensions: Add transport-neutral published services
Socket-exposed APIs previously required extension authors and clients
to use gRPC registration and generated protobuf clients directly.
Add typed service definitions generated from handwritten Go contracts.
The SDK and daemon adapt them to the existing gRPC transport
internally, while clients resolve the same handwritten interface
through an Engine connection. Keep servicegrpcv0 as the raw gRPC
escape hatch.
Signed-off-by: Paweł Gronowski <pawel.gronowski@docker.com>
diff --git daemon/extensions.go daemon/extensions.go
index cb525089d5..7c01785a0e 100644
--- daemon/extensions.go
+++ daemon/extensions.go
@@ -8,11 +8,13 @@
"github.com/containerd/log"
"github.com/moby/moby/v2/daemon/config"
"github.com/moby/moby/v2/daemon/internal/rootless"
+ servicev0 "github.com/moby/moby/v2/extpoints/service/v0"
servicegrpcv0 "github.com/moby/moby/v2/extpoints/servicegrpc/v0"
"github.com/moby/moby/v2/internal/extensions"
"github.com/moby/moby/v2/internal/extensions/grpcproxy"
"github.com/moby/moby/v2/internal/extensions/host"
"github.com/moby/moby/v2/internal/extensions/serverpoint"
+ "github.com/moby/moby/v2/internal/extensions/servicegrpc"
"github.com/moby/moby/v2/pkg/homedir"
"google.golang.org/grpc"
)
@@ -26,7 +28,7 @@ func setupExtensionHost(ctx context.Context, cfg *config.Config) (*host.Host, er
ClientProviders: clientProviders(),
DependencyProviders: dependencyProviders(),
ExtensionConfig: extensionConfig(cfg),
- ExposeOnlyPoints: []extensions.PointID{servicegrpcv0.Point.ID()},
+ ExposeOnlyPoints: []extensions.PointID{servicev0.Point.ID(), servicegrpcv0.Point.ID()},
})
}
@@ -79,9 +81,9 @@ func defaultExtensionDir() (string, error) {
return filepath.Join(libexecDir, "docker", "moby-extensions"), nil
}
-// ExposeExtensionServices publishes services selected through service.grpc on
-// the API socket. In-process services are registered on gs; out-of-process
-// services are proxied by name. Service-name collisions fail startup.
+// ExposeExtensionServices publishes typed and raw extension services on the API
+// socket. In-process services are registered on gs; out-of-process services are
+// proxied by name. Service-name collisions fail startup.
func (daemon *Daemon) ExposeExtensionServices(gs *grpc.Server) (*grpcproxy.Proxy, error) {
if daemon.extensionHost == nil {
return nil, nil
@@ -91,10 +93,21 @@ func (daemon *Daemon) ExposeExtensionServices(gs *grpc.Server) (*grpcproxy.Proxy
reserved[name] = struct{}{}
}
- inproc, err := servicegrpcv0.Collect(daemon.extensionHost)
+ typed, err := servicev0.Collect(daemon.extensionHost)
if err != nil {
return nil, err
}
+ var inproc []servicegrpc.Service
+ for _, registration := range typed {
+ inproc = append(inproc, servicegrpc.Adapt(registration))
+ }
+ raw, err := servicegrpcv0.Collect(daemon.extensionHost)
+ if err != nil {
+ return nil, err
+ }
+ for _, svc := range raw {
+ inproc = append(inproc, servicegrpc.Service{Name: svc.Name, Desc: svc.Desc, Impl: svc.Impl})
+ }
for _, svc := range inproc {
if _, taken := reserved[svc.Name]; taken {
return nil, fmt.Errorf("in-process extension cannot expose gRPC service %q: it is already served", svc.Name)
@@ -106,7 +119,16 @@ func (daemon *Daemon) ExposeExtensionServices(gs *grpc.Server) (*grpcproxy.Proxy
}
var backends []grpcproxy.Backend
- for ext, names := range daemon.extensionHost.ServicesForPoint(servicegrpcv0.Point.ID()) {
+ typedServices := daemon.extensionHost.ServicesForPoint(servicev0.Point.ID())
+ rawServices := daemon.extensionHost.ServicesForPoint(servicegrpcv0.Point.ID())
+ backendServices := make(map[extensions.ExtensionID][]string, len(typedServices)+len(rawServices))
+ for ext, names := range typedServices {
+ backendServices[ext] = append(backendServices[ext], names...)
+ }
+ for ext, names := range rawServices {
+ backendServices[ext] = append(backendServices[ext], names...)
+ }
+ for ext, names := range backendServices {
if conn, ok := daemon.extensionHost.Conn(ext); ok {
backends = append(backends, grpcproxy.Backend{ID: string(ext), Conn: conn, Services: names})
}
diff --git extensions/client/client.go extensions/client/client.go
new file mode 100644
index 0000000000..4481c41e53
--- /dev/null
+++ extensions/client/client.go
@@ -0,0 +1,64 @@
+// Package extensionclient resolves typed extension services published on the
+// Engine API connection.
+package extensionclient
+
+import (
+ "context"
+ "errors"
+ "fmt"
+ "net"
+
+ engineclient "github.com/moby/moby/client"
+ servicev0 "github.com/moby/moby/v2/extpoints/service/v0"
+ "github.com/moby/moby/v2/internal/extensions/servicegrpc"
+ "google.golang.org/grpc"
+ "google.golang.org/grpc/credentials/insecure"
+)
+
+// Client owns one reusable extension RPC connection.
+type Client struct {
+ conn *grpc.ClientConn
+}
+
+// New derives an extension RPC connection from engine's public dialer.
+// It supports local Unix sockets and Windows named pipes.
+func New(engine *engineclient.Client) (*Client, error) {
+ if engine == nil {
+ return nil, errors.New("extension client: engine client is nil")
+ }
+ host, err := engineclient.ParseHostURL(engine.DaemonHost())
+ if err != nil {
+ return nil, fmt.Errorf("extension client: parse engine host: %w", err)
+ }
+ if host.Scheme != "unix" && host.Scheme != "npipe" {
+ return nil, fmt.Errorf("extension client: unsupported engine host scheme %q", host.Scheme)
+ }
+ dialer := engine.Dialer()
+ conn, err := grpc.NewClient("passthrough:///moby-engine",
+ grpc.WithTransportCredentials(insecure.NewCredentials()),
+ grpc.WithContextDialer(func(ctx context.Context, _ string) (net.Conn, error) {
+ return dialer(ctx)
+ }),
+ )
+ if err != nil {
+ return nil, fmt.Errorf("extension client: create connection: %w", err)
+ }
+ return &Client{conn: conn}, nil
+}
+
+// Close closes only the extension RPC connection.
+func (c *Client) Close() error {
+ if c == nil || c.conn == nil {
+ return nil
+ }
+ return c.conn.Close()
+}
+
+// Resolve constructs the handwritten service interface over c's shared
+// connection.
+func Resolve[T any](c *Client, service servicev0.Definition[T]) T {
+ if c == nil || c.conn == nil {
+ panic("extension client: nil client")
+ }
+ return service.NewClient(servicegrpc.NewInvoker(c.conn))
+}
diff --git extensions/client/client_test.go extensions/client/client_test.go
new file mode 100644
index 0000000000..f260f65a38
--- /dev/null
+++ extensions/client/client_test.go
@@ -0,0 +1,69 @@
+package extensionclient
+
+import (
+ "context"
+ "net"
+ "sync/atomic"
+ "testing"
+ "time"
+
+ engineclient "github.com/moby/moby/client"
+ servicev0 "github.com/moby/moby/v2/extpoints/service/v0"
+ "gotest.tools/v3/assert"
+)
+
+func TestNewRejectsNilEngineClient(t *testing.T) {
+ client, err := New(nil)
+ assert.Assert(t, client == nil)
+ assert.ErrorContains(t, err, "engine client is nil")
+}
+
+func TestNewRejectsRemoteEngineHost(t *testing.T) {
+ engine, err := engineclient.New(engineclient.WithHost("tcp://127.0.0.1:2375"))
+ assert.NilError(t, err)
+ defer engine.Close()
+
+ client, err := New(engine)
+ assert.Assert(t, client == nil)
+ assert.ErrorContains(t, err, `unsupported engine host scheme "tcp"`)
+}
+
+func TestNilClientClose(t *testing.T) {
+ var client *Client
+ assert.NilError(t, client.Close())
+}
+
+func TestNewReusesOneConnectionForResolvedServices(t *testing.T) {
+ clientConn, serverConn := net.Pipe()
+ defer serverConn.Close()
+ var dials atomic.Int32
+ dialed := make(chan struct{})
+ engine, err := engineclient.New(engineclient.WithDialContext(func(context.Context, string, string) (net.Conn, error) {
+ if dials.Add(1) == 1 {
+ close(dialed)
+ }
+ return clientConn, nil
+ }))
+ assert.NilError(t, err)
+ defer engine.Close()
+ rpc, err := New(engine)
+ assert.NilError(t, err)
+ defer rpc.Close()
+
+ definition := servicev0.Define("example.Service", func(servicev0.Registrar, serviceMarker) {}, func(servicev0.Invoker) serviceMarker { return serviceMarker{} })
+ Resolve(rpc, definition)
+ Resolve(rpc, definition)
+ assert.Equal(t, dials.Load(), int32(0), "grpc connects lazily")
+
+ ctx, cancel := context.WithTimeout(context.Background(), time.Second)
+ defer cancel()
+ rpc.conn.Connect()
+ select {
+ case <-dialed:
+ case <-ctx.Done():
+ t.Fatal("timed out waiting for connection dial")
+ }
+ assert.Equal(t, dials.Load(), int32(1))
+}
+
+type serviceMarker struct{}
diff --git extpoints/service/v0/service.go extpoints/service/v0/service.go
new file mode 100644
index 0000000000..fb291bb932
--- /dev/null
+++ extpoints/service/v0/service.go
@@ -0,0 +1,113 @@
+// Package servicev0 defines transport-neutral typed services that extensions
+// publish on the daemon API socket.
+package servicev0
+
+import (
+ "context"
+ "fmt"
+
+ "github.com/moby/moby/v2/internal/extensions"
+)
+
+// Point is the infrastructure point used for typed published services.
+var Point = extensions.DefinePoint[Provider]("org.mobyproject.extension.service.v0")
+
+// Invoker performs one unary call on a named service method.
+type Invoker interface {
+ Invoke(ctx context.Context, service, method string, request, response any) error
+}
+
+// UnaryHandler handles one unary request.
+type UnaryHandler func(ctx context.Context, request any) (any, error)
+
+// Method describes one unary service method and constructs its request value.
+type Method struct {
+ Name string
+ NewRequest func() any
+ Handler UnaryHandler
+}
+
+// Registration is a transport-neutral registered service implementation.
+type Registration struct {
+ Service string
+ Methods []Method
+}
+
+// Registrar accepts a transport-neutral service registration.
+type Registrar interface {
+ RegisterService(Registration)
+}
+
+// Definition binds a stable service identity to generated server and client
+// adapters for handwritten interface T.
+type Definition[T any] struct {
+ service string
+ register func(Registrar, T)
+ client func(Invoker) T
+}
+
+// Define creates a typed service definition.
+func Define[T any](service string, register func(Registrar, T), client func(Invoker) T) Definition[T] {
+ if service == "" || register == nil || client == nil {
+ panic("servicev0: incomplete service definition")
+ }
+ return Definition[T]{service: service, register: register, client: client}
+}
+
+// Expose declares impl as a typed service published through Point.
+func (d Definition[T]) Expose(impl T) extensions.Provider {
+ if d.service == "" || d.register == nil || d.client == nil {
+ panic("servicev0: incomplete service definition")
+ }
+ return Point.Provide(provider{
+ register: func(registrar Registrar) {
+ d.register(registrar, impl)
+ },
+ })
+}
+
+// NewClient constructs the handwritten service interface over invoker.
+func (d Definition[T]) NewClient(invoker Invoker) T {
+ if d.service == "" || d.register == nil || d.client == nil {
+ panic("servicev0: incomplete service definition")
+ }
+ if invoker == nil {
+ panic(fmt.Sprintf("servicev0: nil invoker for service %q", d.service))
+ }
+ return d.client(invoker)
+}
+
+// Collect gathers typed service registrations without binding a transport.
+func Collect(resolver extensions.Resolver) ([]Registration, error) {
+ providers, err := Point.All(resolver)
+ if err != nil {
+ return nil, err
+ }
+ var registrations collector
+ for _, provider := range providers {
+ provider.Impl.RegisterServices(®istrations)
+ }
+ return registrations.services, nil
+}
+
+// Provider is a typed service implementation published through Point.
+type Provider interface {
+ // RegisterServices registers the provider's typed services.
+ RegisterServices(Registrar)
+}
+
+type provider struct {
+ register func(Registrar)
+}
+
+func (p provider) RegisterServices(registrar Registrar) {
+ p.register(registrar)
+}
+
+type collector struct {
+ services []Registration
+}
+
+func (c *collector) RegisterService(registration Registration) {
+ c.services = append(c.services, registration)
+}
diff --git extpoints/service/v0/service_test.go extpoints/service/v0/service_test.go
new file mode 100644
index 0000000000..2db891b862
--- /dev/null
+++ extpoints/service/v0/service_test.go
@@ -0,0 +1,94 @@
+package servicev0
+
+import (
+ "context"
+ "testing"
+
+ "github.com/moby/moby/v2/internal/extensions"
+ "gotest.tools/v3/assert"
+ is "gotest.tools/v3/assert/cmp"
+)
+
+type greeting interface {
+ Greet(context.Context, *request) (*response, error)
+}
+
+type request struct{ name string }
+type response struct{ message string }
+
+type greeter struct{}
+
+func (greeter) Greet(_ context.Context, req *request) (*response, error) {
+ return &response{message: "hello " + req.name}, nil
+}
+
+type greetingClient struct{ invoker Invoker }
+
+func (c greetingClient) Greet(ctx context.Context, req *request) (*response, error) {
+ out := new(response)
+ if err := c.invoker.Invoke(ctx, "example.Greeter", "Greet", req, out); err != nil {
+ return nil, err
+ }
+ return out, nil
+}
+
+var greetingService = Define("example.Greeter", func(r Registrar, impl greeting) {
+ r.RegisterService(Registration{Service: "example.Greeter", Methods: []Method{{
+ Name: "Greet", NewRequest: func() any { return new(request) },
+ Handler: func(ctx context.Context, req any) (any, error) {
+ return impl.Greet(ctx, req.(*request))
+ },
+ }}})
+}, func(invoker Invoker) greeting { return greetingClient{invoker: invoker} })
+
+func TestDefinitionExposeAndCollect(t *testing.T) {
+ provider := greetingService.Expose(greeter{})
+ resolver := staticResolver{providers: []extensions.ResolvedProvider{{
+ Extension: "example.greeter.v1",
+ Impl: provider.Impl,
+ }}}
+
+ services, err := Collect(resolver)
+ assert.NilError(t, err)
+ assert.Check(t, is.Len(services, 1))
+ assert.Equal(t, services[0].Service, "example.Greeter")
+ assert.Check(t, is.Len(services[0].Methods, 1))
+ assert.Equal(t, services[0].Methods[0].Name, "Greet")
+
+ got, err := services[0].Methods[0].Handler(context.Background(), &request{name: "world"})
+ assert.NilError(t, err)
+ assert.Equal(t, got.(*response).message, "hello world")
+}
+
+type staticResolver struct {
+ providers []extensions.ResolvedProvider
+}
+
+func (r staticResolver) Provider(extensions.PointID, extensions.ExtensionID) (any, error) {
+ return r.providers[0].Impl, nil
+}
+
+func (r staticResolver) Providers(extensions.PointID) []extensions.ResolvedProvider {
+ return r.providers
+}
+
+func TestDefinitionNewClient(t *testing.T) {
+ invoker := &testInvoker{}
+ client := greetingService.NewClient(invoker)
+ got, err := client.Greet(context.Background(), &request{name: "world"})
+ assert.NilError(t, err)
+ assert.Equal(t, got.message, "hello world")
+ assert.Equal(t, invoker.service, "example.Greeter")
+ assert.Equal(t, invoker.method, "Greet")
+}
+
+type testInvoker struct {
+ service string
+ method string
+}
+
+func (i *testInvoker) Invoke(_ context.Context, service, method string, req, resp any) error {
+ i.service, i.method = service, method
+ resp.(*response).message = "hello " + req.(*request).name
+ return nil
+}
diff --git integration/extension/socket_test.go integration/extension/socket_test.go
index ca3961267d..1a145c0cbc 100644
--- integration/extension/socket_test.go
+++ integration/extension/socket_test.go
@@ -9,17 +9,17 @@
"strconv"
"testing"
+ extensionclient "github.com/moby/moby/v2/extensions/client"
"github.com/moby/moby/v2/integration/extension/testdata/greeter"
+ greeterv0 "github.com/moby/moby/v2/internal/extensions/example/greeter/v0"
greeterpb "github.com/moby/moby/v2/internal/extensions/example/greeter/v0/protogen"
"github.com/moby/moby/v2/internal/testutil"
"github.com/moby/moby/v2/internal/testutil/daemon"
- "google.golang.org/grpc"
- "google.golang.org/grpc/credentials/insecure"
"gotest.tools/v3/assert"
"gotest.tools/v3/skip"
)
-func TestSocketExposedGRPCService(t *testing.T) {
+func TestSocketExposedTypedService(t *testing.T) {
skip.If(t, testEnv.IsRemoteDaemon, "the extension binary must be on the daemon's host")
ctx := testutil.StartSpan(baseContext, t)
@@ -34,13 +34,15 @@ func TestSocketExposedGRPCService(t *testing.T) {
d.Start(t, startArgs...)
defer d.Stop(t)
- conn, err := grpc.NewClient(d.Sock(), grpc.WithTransportCredentials(insecure.NewCredentials()))
+ engine := d.NewClientT(t)
+ extClient, err := extensionclient.New(engine)
assert.NilError(t, err)
- defer conn.Close()
+ defer extClient.Close()
- resp, err := greeterpb.NewGreeterClient(conn).Greet(ctx, &greeterpb.HelloRequest{Name: "world"})
+ client := extensionclient.Resolve(extClient, greeterpb.Service)
+ resp, err := client.Greet(ctx, &greeterv0.HelloRequest{Name: "world"})
assert.NilError(t, err)
- assert.Equal(t, resp.GetMessage(), "hello world")
+ assert.Equal(t, resp.Message, "hello world")
}
func buildGreeterExtension(ctx context.Context, t *testing.T) string {
diff --git integration/extension/testdata/greeter/cmd/greeter/main.go integration/extension/testdata/greeter/cmd/greeter/main.go
index 99c863911f..dfea57a22e 100644
--- integration/extension/testdata/greeter/cmd/greeter/main.go
+++ integration/extension/testdata/greeter/cmd/greeter/main.go
@@ -3,11 +3,10 @@
package main
import (
- servicegrpcv0 "github.com/moby/moby/v2/extpoints/servicegrpc/v0"
"github.com/moby/moby/v2/integration/extension/testdata/greeter"
"github.com/moby/moby/v2/internal/extensions/sdk"
)
func main() {
- sdk.Main(greeter.Extension, servicegrpcv0.ServerPoint)
+ sdk.Main(greeter.Extension)
}
diff --git integration/extension/testdata/greeter/greeter.go integration/extension/testdata/greeter/greeter.go
index 6bf5eafaeb..88dff012f2 100644
--- integration/extension/testdata/greeter/greeter.go
+++ integration/extension/testdata/greeter/greeter.go
@@ -4,11 +4,9 @@
import (
"context"
- servicegrpcv0 "github.com/moby/moby/v2/extpoints/servicegrpc/v0"
"github.com/moby/moby/v2/internal/extensions"
greeterv0 "github.com/moby/moby/v2/internal/extensions/example/greeter/v0"
greeterpb "github.com/moby/moby/v2/internal/extensions/example/greeter/v0/protogen"
- "google.golang.org/grpc"
)
// ID is the extension id and binary name.
@@ -20,15 +18,8 @@ func (greeter) Greet(_ context.Context, req *greeterv0.HelloRequest) (*greeterv0
return &greeterv0.HelloReply{Message: "hello " + req.Name}, nil
}
-// expose registers the greeter gRPC service for socket exposure.
-type expose struct{}
-
-func (expose) RegisterServices(r grpc.ServiceRegistrar) {
- greeterpb.ServerPoint.Register(r, greeter{})
-}
-
-// Extension implements the service.grpc point.
+// Extension publishes the typed Greeter service.
var Extension = extensions.New(extensions.Declaration{
ID: ID,
- Providers: []extensions.Provider{servicegrpcv0.Point.Provide(expose{})},
+ Providers: []extensions.Provider{greeterpb.Service.Expose(greeter{})},
})
diff --git integration/extension/testdata/greeterdep/cmd/greeterdep/main.go integration/extension/testdata/greeterdep/cmd/greeterdep/main.go
index 3c5d7d4c4a..d109e5fde8 100644
--- integration/extension/testdata/greeterdep/cmd/greeterdep/main.go
+++ integration/extension/testdata/greeterdep/cmd/greeterdep/main.go
@@ -10,7 +10,7 @@
"syscall"
"github.com/moby/moby/v2/integration/extension/testdata/greeterdep"
- greeterpb "github.com/moby/moby/v2/internal/extensions/example/greeter/v0/protogen"
+ greeterpointpb "github.com/moby/moby/v2/internal/extensions/example/greeterpoint/v0/protogen"
"github.com/moby/moby/v2/internal/extensions/sdk"
)
@@ -23,7 +23,7 @@ func main() {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
- srv.Depends(greeterpb.ClientPoint)
+ srv.Depends(greeterpointpb.ClientPoint)
if err := srv.Listen(ctx); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
diff --git integration/extension/testdata/greeterdep/greeterdep.go integration/extension/testdata/greeterdep/greeterdep.go
index d029b1eaff..3c70409c95 100644
--- integration/extension/testdata/greeterdep/greeterdep.go
+++ integration/extension/testdata/greeterdep/greeterdep.go
@@ -6,7 +6,7 @@
"fmt"
"github.com/moby/moby/v2/internal/extensions"
- greeterv0 "github.com/moby/moby/v2/internal/extensions/example/greeter/v0"
+ greeterpointv0 "github.com/moby/moby/v2/internal/extensions/example/greeterpoint/v0"
)
// ID is the extension id and binary name.
@@ -14,7 +14,7 @@
// initialize calls the greeter dependency and verifies the reply.
func initialize(ctx context.Context, _ extensions.Config, r extensions.Resolver) error {
- reply, err := greeterv0.Greet(ctx, r, &greeterv0.HelloRequest{Name: "dep"})
+ reply, err := greeterpointv0.Greet(ctx, r, &greeterpointv0.HelloRequest{Name: "dep"})
if err != nil {
return fmt.Errorf("greeterdep: call greeter dependency: %w", err)
}
@@ -27,6 +27,6 @@ func initialize(ctx context.Context, _ extensions.Config, r extensions.Resolver)
// Extension declares and uses a dependency on the greeter point.
var Extension = extensions.New(extensions.Declaration{
ID: ID,
- Dependencies: []extensions.Dependency{greeterv0.Point.Dependency()},
+ Dependencies: []extensions.Dependency{greeterpointv0.Point.Dependency()},
Init: initialize,
})
diff --git internal/extensions/cmd/mobyextgen/main.go internal/extensions/cmd/mobyextgen/main.go
index 27777dbf68..2ac34ca9d0 100644
--- internal/extensions/cmd/mobyextgen/main.go
+++ internal/extensions/cmd/mobyextgen/main.go
@@ -1,5 +1,5 @@
-// Command mobyextgen generates an extension point's wire contract and transport
-// code from a Go interface and pb-tagged message structs.
+// Command mobyextgen generates an extension point or service's wire contract
+// and transport code from a Go interface and pb-tagged message structs.
//
// Run with no arguments from the point's package, typically through `go generate`.
//
@@ -135,6 +135,7 @@ type point struct {
iface string // Go service interface name
isPoint bool // whether the contract declares an extensions.Point
isSingle bool // whether the point was declared with DefineSinglePoint
+ published bool // whether a non-point service is published on the API socket
methods []method
messages []message
}
@@ -212,7 +213,11 @@ func parsePoint(dir string) (point, error) {
pt.service = svc.service
if svc.iface != "" {
pt.iface, pt.id = svc.iface, svc.pkg
+ pt.published = svc.published
} else {
+ if svc.published {
+ return point{}, errors.New("mobyextgen:publish is only valid on a non-point service interface")
+ }
iface, id, single, err := findDefinePoint(files)
if err != nil {
return point{}, err
@@ -282,8 +287,9 @@ func findDefinePoint(files []*ast.File) (iface, id string, single bool, err erro
return iface, id, single, nil
}
-// servicePragma declares a contract's gRPC service name, and reservedPragma
-// preserves a removed field number.
+// servicePragma declares a contract's wire service name, publishPragma opts a
+// non-point interface into typed publication, and reservedPragma preserves a
+// removed field number.
//
// //mobyextgen:service=CreateSpecHook
// var Point = extensions.DefinePoint[Hook]("...create_spec.v0")
@@ -292,15 +298,17 @@ func findDefinePoint(files []*ast.File) (iface, id string, single bool, err erro
// type PointDeclaration struct{ ... }
const (
servicePragma = "//mobyextgen:service="
+ publishPragma = "//mobyextgen:publish"
reservedPragma = "//mobyextgen:reserved="
)
// serviceDecl is the service pragma result. Point mode supplies only service;
// service mode supplies the interface and fully-qualified proto package too.
type serviceDecl struct {
- service string
- iface string
- pkg string
+ service string
+ iface string
+ pkg string
+ published bool
}
// findServicePragma returns the declared gRPC service. The name is required in
@@ -318,6 +326,9 @@ func findServicePragma(files []*ast.File) (serviceDecl, error) {
for _, spec := range gd.Specs {
values := declPragmas(gd, spec, servicePragma)
if len(values) == 0 {
+ if len(declPragmas(gd, spec, publishPragma)) > 0 {
+ return serviceDecl{}, errors.New("mobyextgen:publish requires a mobyextgen:service pragma on the same interface")
+ }
continue
}
for _, value := range values {
@@ -330,17 +341,18 @@ func findServicePragma(files []*ast.File) (serviceDecl, error) {
ts, isType := spec.(*ast.TypeSpec)
if !isType {
- found = serviceDecl{service: value}
+ found = serviceDecl{service: value, published: len(declPragmas(gd, spec, publishPragma)) > 0}
continue
}
if _, isIface := ts.Type.(*ast.InterfaceType); !isIface {
return serviceDecl{}, fmt.Errorf("%s on %s: the pragma may only document an interface or the point", servicePragma, ts.Name.Name)
}
+ published := len(declPragmas(gd, spec, publishPragma)) > 0
pkg, service, ok := strings.Cut(reverse(value), ".")
if !ok {
return serviceDecl{}, fmt.Errorf("%s on interface %s: want a fully-qualified name like my.proto.package.%s", servicePragma, ts.Name.Name, ts.Name.Name)
}
- found = serviceDecl{service: reverse(pkg), iface: ts.Name.Name, pkg: reverse(service)}
+ found = serviceDecl{service: reverse(pkg), iface: ts.Name.Name, pkg: reverse(service), published: published}
}
}
}
@@ -696,14 +708,19 @@ func emitWire(pt point) ([]byte, error) {
fmt.Fprintln(&b, ` clientpoint "github.com/moby/moby/v2/internal/extensions/clientpoint"`)
fmt.Fprintln(&b, ` serverpoint "github.com/moby/moby/v2/internal/extensions/serverpoint"`)
}
+ if pt.published {
+ fmt.Fprintln(&b, ` servicev0 "github.com/moby/moby/v2/extpoints/service/v0"`)
+ }
fmt.Fprintf(&b, " %s %q\n", cpkg, pt.importPath)
fmt.Fprintln(&b, ` grpc "google.golang.org/grpc"`)
fmt.Fprintln(&b, ")")
- fmt.Fprintf(&b, `
-// serviceName is the point's fully-qualified gRPC service name.
-const serviceName = %q
-`, pt.grpcService())
+ if pt.published {
+ fmt.Fprintln(&b, "\n// serviceName is the contract's fully-qualified gRPC service name.")
+ } else {
+ fmt.Fprintln(&b, "\n// serviceName is the point's fully-qualified gRPC service name.")
+ }
+ fmt.Fprintf(&b, "const serviceName = %q\n", pt.grpcService())
fmt.Fprintln(&b, "\nconst (")
for _, m := range pt.methods {
@@ -711,20 +728,32 @@ func emitWire(pt point) ([]byte, error) {
}
fmt.Fprintln(&b, ")")
- fmt.Fprintf(&b, `
+ if pt.published {
+ fmt.Fprintf(&b, `
+// %[1]sServer is the server side of the contract's gRPC service. It is the
+// proto-level shape of the service, not the handwritten Go interface: a contract
+// method returning a bare error returns an empty response message here.
+type %[1]sServer interface {
+`, svc)
+ } else {
+ fmt.Fprintf(&b, `
// %[1]sServer is the server side of the point's gRPC service. It is the
// proto-level shape of the point, not the point's Go interface: a contract
// method returning a bare error returns an empty response message here.
type %[1]sServer interface {
`, svc)
+ }
for _, m := range pt.methods {
fmt.Fprintf(&b, "\t%s(context.Context, *%s) (*%s, error)\n", m.name, m.request, m.response)
}
fmt.Fprintln(&b, "}")
- fmt.Fprintf(&b, `
-// serviceDesc describes the point's gRPC service to a server. HandlerType is
-// what a registrar type-checks an implementation against, so registering the
+ if pt.published {
+ fmt.Fprintln(&b, "\n// serviceDesc describes the contract's gRPC service to a server. HandlerType is")
+ } else {
+ fmt.Fprintln(&b, "\n// serviceDesc describes the point's gRPC service to a server. HandlerType is")
+ }
+ fmt.Fprintf(&b, `// what a registrar type-checks an implementation against, so registering the
// wrong provider for this point is caught at registration.
var serviceDesc = grpc.ServiceDesc{
ServiceName: serviceName,
@@ -755,20 +784,32 @@ func %[1]s(srv any, ctx context.Context, dec func(any) error, interceptor grpc.U
`, handlerName(m), m.request, svc, m.name, methodConst(m))
}
- fmt.Fprintf(&b, `
+ if pt.published {
+ fmt.Fprintf(&b, `
+// %[1]sClient calls the contract's raw gRPC service. It is exported so a client
+// outside the framework - one calling a service an extension publishes on the
+// API socket - can reach it with a plain gRPC client.
+type %[1]sClient interface {
+`, svc)
+ } else {
+ fmt.Fprintf(&b, `
// %[1]sClient calls the point's gRPC service. It is exported so a client
// outside the framework - one calling a service an extension publishes on the
// API socket - can reach it with a plain gRPC client.
type %[1]sClient interface {
`, svc)
+ }
for _, m := range pt.methods {
fmt.Fprintf(&b, "\t%s(ctx context.Context, in *%s, opts ...grpc.CallOption) (*%s, error)\n", m.name, m.request, m.response)
}
fmt.Fprintln(&b, "}")
- fmt.Fprintf(&b, `
-// New%[1]sClient returns a client for the point's gRPC service on cc.
-func New%[1]sClient(cc grpc.ClientConnInterface) %[1]sClient { return &serviceClient{cc: cc} }
+ if pt.published {
+ fmt.Fprintf(&b, "\n// New%[1]sClient returns a client for the contract's raw gRPC service on cc.\n", svc)
+ } else {
+ fmt.Fprintf(&b, "\n// New%[1]sClient returns a client for the point's gRPC service on cc.\n", svc)
+ }
+ fmt.Fprintf(&b, `func New%[1]sClient(cc grpc.ClientConnInterface) %[1]sClient { return &serviceClient{cc: cc} }
type serviceClient struct{ cc grpc.ClientConnInterface }
`, svc)
@@ -817,6 +858,28 @@ func NewClient(conn grpc.ClientConnInterface) %[2]s.%[3]s {
return &grpcClient{client: New%[1]sClient(conn)}
}
`, svc, cpkg, iface)
+ if pt.published {
+ fmt.Fprintf(&b, `
+
+// Service is the transport-neutral typed definition of the published %[1]s service.
+var Service = servicev0.Define(serviceName, func(r servicev0.Registrar, impl %[2]s.%[3]s) {
+ r.RegisterService(servicev0.Registration{
+ Service: serviceName,
+ Methods: []servicev0.Method{
+`, svc, cpkg, iface)
+ for _, m := range pt.methods {
+ fmt.Fprintf(&b, ` {Name: %q, NewRequest: func() any { return new(%s) }, Handler: func(ctx context.Context, request any) (any, error) {
+ return (&grpcServer{impl: impl}).%s(ctx, request.(*%s))
+ }},
+`, m.name, m.request, m.name, m.request)
+ }
+ fmt.Fprintf(&b, ` },
+ })
+ }, func(invoker servicev0.Invoker) %[1]s.%[2]s {
+ return &typedClient{invoker: invoker}
+ })
+`, cpkg, iface)
+ }
}
fmt.Fprintf(&b, `
@@ -854,6 +917,13 @@ type grpcClient struct {
client %sClient
}
`, svc)
+ if pt.published {
+ fmt.Fprint(&b, `
+type typedClient struct {
+ invoker servicev0.Invoker
+}
+`)
+ }
for _, m := range pt.methods {
if m.bareError {
@@ -874,6 +944,25 @@ func (c *grpcClient) %[1]s(ctx context.Context, req *%[2]s.%[3]s) (*%[2]s.%[5]s,
}
`, m.name, cpkg, m.request, lowerFirst(m.request), m.response, lowerFirst(m.response))
}
+ if pt.published {
+ if m.bareError {
+ fmt.Fprintf(&b, `
+func (c *typedClient) %[1]s(ctx context.Context, req *%[2]s.%[3]s) error {
+ return c.invoker.Invoke(ctx, serviceName, %[4]q, %[5]sToProto(req), new(%[6]s))
+}
+`, m.name, cpkg, m.request, m.name, lowerFirst(m.request), m.response)
+ } else {
+ fmt.Fprintf(&b, `
+func (c *typedClient) %[1]s(ctx context.Context, req *%[2]s.%[3]s) (*%[2]s.%[4]s, error) {
+ out := new(%[4]s)
+ if err := c.invoker.Invoke(ctx, serviceName, %[5]q, %[6]sToProto(req), out); err != nil {
+ return nil, err
+ }
+ return %[7]sFromProto(out), nil
+}
+`, m.name, cpkg, m.request, m.response, m.name, lowerFirst(m.request), lowerFirst(m.response))
+ }
+ }
}
for _, msg := range pt.messages {
diff --git internal/extensions/cmd/mobyextgen/main_test.go internal/extensions/cmd/mobyextgen/main_test.go
index ead639c8d1..374936a9d4 100644
--- internal/extensions/cmd/mobyextgen/main_test.go
+++ internal/extensions/cmd/mobyextgen/main_test.go
@@ -259,6 +259,48 @@ type Runtime interface{ Do(ctx interface{
assert.Check(t, !strings.Contains(src, "clientpoint"), "a non-point contract must not import the point packages:\n%s", src)
}
+func TestPublishedServiceContract(t *testing.T) {
+ const contract = `package p
+type Req struct{ Name string ` + "`pb:\"1\"`" + ` }
+type Resp struct{ Ok bool ` + "`pb:\"1\"`" + ` }
+
+//mobyextgen:publish
+//mobyextgen:service=my.proto.pkg.v1.Runtime
+type Runtime interface{ Do(ctx interface{}, req *Req) (*Resp, error) }
+`
+ pt, err := parseSource(t, contract)
+ assert.NilError(t, err)
+ assert.Check(t, pt.published)
+ pt.importPath = "example.com/p"
+ wire, err := emitWire(pt)
+ assert.NilError(t, err)
+ src := string(wire)
+ assert.Check(t, strings.Contains(src, "var Service = servicev0.Define(serviceName"), src)
+ assert.Check(t, strings.Contains(src, `Invoke(ctx, serviceName, "Do"`), src)
+}
+
+func TestPublishRejectedOnPoint(t *testing.T) {
+ _, err := parseSource(t, `package p
+import "github.com/moby/moby/v2/internal/extensions"
+type Req struct{}
+type Resp struct{}
+type Runtime interface{ Do(ctx interface{}, req *Req) (*Resp, error) }
+//mobyextgen:publish
+//mobyextgen:service=Runtime
+var Point = extensions.DefinePoint[Runtime]("example.point.v1")
+`)
+ assert.ErrorContains(t, err, "publish is only valid on a non-point service interface")
+}
+
+func TestPublishRequiresServicePragma(t *testing.T) {
+ _, err := parseSource(t, `package p
+//mobyextgen:publish
+type Runtime interface{
+}
+`)
+ assert.ErrorContains(t, err, "publish requires a mobyextgen:service pragma")
+}
+
func TestServicePragmaOnANonInterface(t *testing.T) {
_, err := parseSource(t, `package p
type Req struct{ Name string `+"`pb:\"1\"`"+` }
diff --git internal/extensions/docs/AUTHORING.md internal/extensions/docs/AUTHORING.md
index 8f889522e2..58e1a20d2c 100644
--- internal/extensions/docs/AUTHORING.md
+++ internal/extensions/docs/AUTHORING.md
@@ -189,43 +189,59 @@ ### 5. Support separate-binary providers
This list is the boundary for launched providers: an unlisted declared point is rejected, while any installed extension may provide a listed point.
See [DESIGN.md#discovery-security](./DESIGN.md#discovery-security).
-An externally published gRPC service is not added to this list.
+An externally published service is not added to this list.
It uses socket exposure instead.
### Socket-exposed services
-Socket exposure is an extension's own gRPC API, not a daemon-called point.
-The daemon forwards raw calls by service name without importing the proto.
-Opt in with `service.grpc` and register through the supplied registrar:
+Socket exposure publishes an extension's own typed API, not a daemon-called
+point.
+Write the same handwritten interface and messages as for a point, but put the
+service pragma and the explicit publication marker on the interface:
```go
-type expose struct{}
-
-func (expose) RegisterServices(r grpc.ServiceRegistrar) {
- mypb.RegisterMyServiceServer(r, impl) // or mypb.ServerPoint.Register(r, impl)
+//mobyextgen:publish
+//mobyextgen:service=com.example.myext.v1.Greeter
+type Greeter interface {
+ Greet(context.Context, *HelloRequest) (*HelloReply, error)
}
+```
+
+Generation emits `mypb.Service`, a transport-neutral
+`servicev0.Definition[Greeter]`.
+Publish an ordinary Go implementation without importing gRPC or protobuf types:
+```go
var Extension = extensions.New(extensions.Declaration{
ID: "com.example.myext.v1",
- Providers: []extensions.Provider{servicegrpcv0.Point.Provide(expose{})},
+ Providers: []extensions.Provider{mypb.Service.Expose(greeter{})},
})
```
-The same implementation works in both modes: in-process, the daemon supplies its gRPC server; out-of-process, the SDK supplies the extension server, records names, and the daemon proxies matching calls.
+The SDK knows the typed-service infrastructure registration, so a binary that
+only publishes typed services uses `sdk.Main(Extension)`.
+Continue to pass generated `ServerPoint` values explicitly for ordinary points.
-An out-of-process binary registers `ServerPoint` for every point it provides:
+Clients derive one reusable extension RPC connection from an existing Engine
+client and resolve the same handwritten interface:
```go
-srv := sdk.NewServer()
-srv.Register(ext,
- servicegrpcv0.ServerPoint,
- mypointpb.ServerPoint, // include every other point this extension provides
-)
-srv.Listen(ctx)
+rpc, err := extensionclient.New(engine)
+if err != nil { /* handle */ }
+defer rpc.Close()
+greeter := extensionclient.Resolve(rpc, mypb.Service)
+reply, err := greeter.Greet(ctx, &myv1.HelloRequest{Name: "world"})
```
+The public author and client APIs are transport-neutral; the current backend is
+gRPC and the daemon proxies out-of-process services by their stable service
+identity.
+The client currently supports local `unix` and `npipe` Engine hosts only;
+TCP hosts cannot safely recover the Engine connection's TLS and ALPN settings.