This is an automated email from the ASF dual-hosted git repository.

shunping pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new f612c6a13cd Port Wait.On to the Go SDK (#39937)
f612c6a13cd is described below

commit f612c6a13cdf0ad759648ebca51a771e3be91b53
Author: Valient Gough <[email protected]>
AuthorDate: Thu Oct 1 18:03:54 2026 -0700

    Port Wait.On to the Go SDK (#39937)
    
    * Port Wait.On to the Go SDK
    
    * update exclusion comments
---
 CHANGES.md                                       |   1 +
 sdks/go/pkg/beam/transforms/wait/wait.go         | 172 +++++++++++++
 sdks/go/pkg/beam/transforms/wait/wait_test.go    | 223 ++++++++++++++++
 sdks/go/test/integration/integration.go          |  13 +
 sdks/go/test/integration/primitives/wait.go      | 307 +++++++++++++++++++++++
 sdks/go/test/integration/primitives/wait_test.go |  73 ++++++
 6 files changed, 789 insertions(+)

diff --git a/CHANGES.md b/CHANGES.md
index 67db578fc27..48f4a022538 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -68,6 +68,7 @@
 
 ## New Features / Improvements
 
+* (Go) Added `wait.On`, which delays each input window until the corresponding 
windows in its signal PCollections have closed 
([#39909](https://github.com/apache/beam/issues/39909)).
 * (Python) Expanded the SDK worker heap dump 
(`--experiments=enable_heap_dump`) with process RSS, CPython allocator/GC 
stats, and glibc `mallinfo2` native-heap/fragmentation stats to help 
distinguish native-heap from Python-object memory growth 
([#39244](https://github.com/apache/beam/issues/39244)).
 * The `disableCounterMetrics`, `disableStringSetMetrics` and 
`disableBoundedTrieMetrics` experiments are now honored by the Python SDK, as 
they already were in Java (Python) 
([#38746](https://github.com/apache/beam/issues/38746)).
 
diff --git a/sdks/go/pkg/beam/transforms/wait/wait.go 
b/sdks/go/pkg/beam/transforms/wait/wait.go
new file mode 100644
index 00000000000..93f7d00f96a
--- /dev/null
+++ b/sdks/go/pkg/beam/transforms/wait/wait.go
@@ -0,0 +1,172 @@
+// Licensed to the Apache Software Foundation (ASF) under one or more
+// contributor license agreements.  See the NOTICE file distributed with
+// this work for additional information regarding copyright ownership.
+// The ASF licenses this file to You under the Apache License, Version 2.0
+// (the "License"); you may not use this file except in compliance with
+// the License.  You may obtain a copy of the License at
+//
+//    http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+// Package wait delays a PCollection until one or more signal PCollections are
+// ready. It is the Go equivalent of Java's Wait.on.
+//
+// On preserves the main input but holds each of its windows until the mapped
+// window in every signal has closed. Global signals map to the global window;
+// fixed and sliding signals map to the earliest window containing the main
+// window's maximum timestamp. A signal window closes after its watermark
+// passes the window end plus allowed lateness.
+//
+// For example, to finish writing each window to one database before writing it
+// to another:
+//
+//     firstWriteResults := beam.ParDo(s, writeToFirstDB, data)
+//     delayed := wait.On(s, data, firstWriteResults)
+//     beam.ParDo0(s, writeToSecondDB, delayed)
+//
+// A bounded global signal holds every main window until the entire signal is
+// complete. An unbounded global signal never becomes ready. Large allowed
+// lateness values delay readiness by the same amount.
+//
+// Go side-input window restrictions apply. Signal PCollections cannot use
+// session windows, and a global main input cannot wait on a non-global signal.
+// Ordinary and KV PCollections are supported. With no signals, On returns any
+// valid input directly; otherwise, CoGBK PCollections are not supported.
+package wait
+
+import (
+       "fmt"
+       "time"
+
+       "github.com/apache/beam/sdks/v2/go/pkg/beam"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph/window"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph/window/trigger"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/core/typex"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/register"
+)
+
+func init() {
+       register.DoFn3x0[typex.Window, beam.T, func(int)](&collectWindowsFn{})
+       register.Function2x1(keepOneFn)
+       register.Function2x1(passThroughFn)
+       register.Function3x2(passThroughKVFn)
+       register.Emitter1[int]()
+       register.Iter1[int]()
+}
+
+// On returns col unchanged — same elements, type, coder and windowing — but
+// delays each main window until its mapped window has closed in every signal.
+//
+// With no signals, On validates s and col and returns col directly. Otherwise,
+// it panics at pipeline construction time for an invalid scope or PCollection,
+// CoGBK inputs, session-windowed signals, or a non-global signal on a globally
+// windowed main input.
+func On(s beam.Scope, col beam.PCollection, signals ...beam.PCollection) 
beam.PCollection {
+       if !s.IsValid() {
+               panic("wait.On: invalid scope")
+       }
+       if !col.IsValid() {
+               panic("wait.On: invalid input pcollection")
+       }
+       if len(signals) == 0 {
+               return col
+       }
+       if typex.IsCoGBK(col.Type()) {
+               panic(fmt.Sprintf("wait.On: input pcollection must not be a 
CoGBK: %v", col))
+       }
+       mainIsGlobal := col.WindowingStrategy().Fn.Kind == window.GlobalWindows
+       for i, sig := range signals {
+               if !sig.IsValid() {
+                       panic(fmt.Sprintf("wait.On: invalid signal pcollection: 
index %d", i))
+               }
+               if typex.IsCoGBK(sig.Type()) {
+                       panic(fmt.Sprintf("wait.On: signal pcollection must not 
be a CoGBK: index %d: %v", i, sig))
+               }
+               if sig.WindowingStrategy().Fn.Kind == window.Sessions {
+                       panic(fmt.Sprintf("wait.On: signal pcollection must not 
use session windowing (side inputs cannot map session windows): index %d: %v", 
i, sig))
+               }
+               if mainIsGlobal && sig.WindowingStrategy().Fn.Kind != 
window.GlobalWindows {
+                       panic(fmt.Sprintf("wait.On: signal pcollection must be 
globally windowed when the input pcollection is (a global main window cannot be 
mapped to a non-global side-input window): index %d: %v", i, sig))
+               }
+       }
+       s = s.Scope("wait.On")
+       out := col
+       for i, sig := range signals {
+               out = onOne(s.Scope(fmt.Sprintf("Signal(%d)", i)), out, sig)
+       }
+       return out
+}
+
+// onOne delays col on a single signal.
+//
+// The Never trigger delays each nonempty signal window's marker until
+// expiration. The marker is used as a side input to an identity ParDo; its
+// value is not read, but its readiness gates the main input. Empty signal
+// windows become ready when the signal watermark passes them.
+//
+// Marker generation emits at most one marker per signal window in each bundle;
+// the Combine collapses markers across bundles.
+func onOne(s beam.Scope, col, signal beam.PCollection) beam.PCollection {
+       ws := signal.WindowingStrategy()
+       if typex.IsKV(signal.Type()) {
+               signal = beam.DropKey(s, signal)
+       }
+       closed := beam.WindowInto(s, ws.Fn, signal,
+               beam.Trigger(trigger.Never()),
+               beam.PanesDiscard(),
+               
beam.AllowedLateness(time.Duration(ws.AllowedLateness)*time.Millisecond),
+       )
+       markers := beam.ParDo(s, &collectWindowsFn{}, closed)
+       marker := beam.Combine(s, keepOneFn, markers)
+       var out beam.PCollection
+       if typex.IsKV(col.Type()) {
+               out = beam.ParDo(s, passThroughKVFn, col, beam.SideInput{Input: 
marker})
+       } else {
+               out = beam.ParDo(s, passThroughFn, col, beam.SideInput{Input: 
marker})
+       }
+       // ParDo infers a fresh coder for its output. The output is col 
verbatim, so
+       // keep col's coder — as Flatten does — rather than replace one the 
user set.
+       if err := out.SetCoder(col.Coder()); err != nil {
+               panic(fmt.Sprintf("wait.On: cannot preserve the input coder on 
the output: %v", err))
+       }
+       return out
+}
+
+// collectWindowsFn emits one marker per signal window in each bundle. 
Observing
+// the window makes the harness process multi-window elements once per window;
+// ProcessElement emitters preserve that window on their output.
+type collectWindowsFn struct {
+       seen map[typex.Window]struct{}
+}
+
+func (fn *collectWindowsFn) StartBundle(_ func(int)) {
+       fn.seen = make(map[typex.Window]struct{})
+}
+
+func (fn *collectWindowsFn) ProcessElement(w typex.Window, _ beam.T, emit 
func(int)) {
+       if _, ok := fn.seen[w]; ok {
+               return
+       }
+       fn.seen[w] = struct{}{}
+       emit(1)
+}
+
+// keepOneFn collapses markers to one per window.
+func keepOneFn(_, _ int) int {
+       return 1
+}
+
+// passThroughFn returns its input after the marker side input is ready.
+func passThroughFn(elm beam.T, _ func(*int) bool) beam.T {
+       return elm
+}
+
+// passThroughKVFn is passThroughFn for a KV main input.
+func passThroughKVFn(k beam.X, v beam.Y, _ func(*int) bool) (beam.X, beam.Y) {
+       return k, v
+}
diff --git a/sdks/go/pkg/beam/transforms/wait/wait_test.go 
b/sdks/go/pkg/beam/transforms/wait/wait_test.go
new file mode 100644
index 00000000000..8e3fc4245c2
--- /dev/null
+++ b/sdks/go/pkg/beam/transforms/wait/wait_test.go
@@ -0,0 +1,223 @@
+// Licensed to the Apache Software Foundation (ASF) under one or more
+// contributor license agreements.  See the NOTICE file distributed with
+// this work for additional information regarding copyright ownership.
+// The ASF licenses this file to You under the Apache License, Version 2.0
+// (the "License"); you may not use this file except in compliance with
+// the License.  You may obtain a copy of the License at
+//
+//    http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package wait_test
+
+import (
+       "fmt"
+       "strings"
+       "testing"
+       "time"
+
+       "github.com/apache/beam/sdks/v2/go/pkg/beam"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph/window"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/register"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/testing/passert"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/testing/ptest"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/transforms/wait"
+)
+
+func init() {
+       register.Function1x2(waitTestKVFn)
+       register.Function2x1(waitTestFormatKVFn)
+}
+
+// TestMain invokes ptest.Main so the end-to-end tests below run on the
+// configured runner (Prism by default).
+func TestMain(m *testing.M) {
+       ptest.Main(m)
+}
+
+// expectPanic runs f and fails the test unless f panics with a message
+// containing want.
+func expectPanic(t *testing.T, want string, f func()) {
+       t.Helper()
+       defer func() {
+               r := recover()
+               if r == nil {
+                       t.Fatalf("expected panic containing %q, got no panic", 
want)
+               }
+               if msg := fmt.Sprint(r); !strings.Contains(msg, want) {
+                       t.Fatalf("panic message %q does not contain %q", msg, 
want)
+               }
+       }()
+       f()
+}
+
+func TestOn_NoSignals(t *testing.T) {
+       _, s := beam.NewPipelineWithRoot()
+       col := beam.Create(s, 1, 2, 3)
+       out := wait.On(s, col)
+       if out != col {
+               t.Errorf("wait.On with no signals returned %v, want the input 
PCollection %v", out, col)
+       }
+}
+
+func TestOn_NoSignalsSkipsShapeChecks(t *testing.T) {
+       _, s := beam.NewPipelineWithRoot()
+       cogbk := beam.CoGroupByKey(s, kvCol(s, 1, 2, 3), kvCol(s, 4, 5, 6))
+       if out := wait.On(s, cogbk); out != cogbk {
+               t.Errorf("wait.On with no signals on a CoGBK input returned %v, 
want the input PCollection %v", out, cogbk)
+       }
+}
+
+func TestOn_InvalidInputsPanic(t *testing.T) {
+       _, s := beam.NewPipelineWithRoot()
+       col := beam.Create(s, 1, 2, 3)
+       sig := beam.Create(s, "ready")
+
+       expectPanic(t, "wait.On: invalid scope", func() {
+               wait.On(beam.Scope{}, col, sig)
+       })
+       expectPanic(t, "wait.On: invalid input pcollection", func() {
+               wait.On(s, beam.PCollection{}, sig)
+       })
+       expectPanic(t, "wait.On: invalid signal pcollection: index 1", func() {
+               wait.On(s, col, sig, beam.PCollection{})
+       })
+}
+
+func TestOn_SessionSignalPanics(t *testing.T) {
+       _, s := beam.NewPipelineWithRoot()
+       col := beam.Create(s, 1, 2, 3)
+       sig := beam.WindowInto(s, window.NewSessions(time.Minute), 
beam.Create(s, "ready"))
+
+       expectPanic(t, "wait.On: signal pcollection must not use session 
windowing (side inputs cannot map session windows): index 0", func() {
+               wait.On(s, col, sig)
+       })
+}
+
+func TestOn_Identity(t *testing.T) {
+       ptest.BuildAndRun(t, func(s beam.Scope) {
+               col := beam.Create(s, 1, 2, 3)
+               sig := beam.Create(s, "ready")
+               out := wait.On(s, col, sig)
+               passert.Equals(s, out, 1, 2, 3)
+       })
+}
+
+func TestOn_MultipleSignals(t *testing.T) {
+       ptest.BuildAndRun(t, func(s beam.Scope) {
+               col := beam.Create(s, 1, 2, 3)
+               sigA := beam.Create(s, "a")
+               sigB := beam.Create(s, 1.5)
+               passert.Equals(s, wait.On(s, col, sigA, sigB), 1, 2, 3)
+       })
+}
+
+func TestOn_GlobalMainNonGlobalSignalPanics(t *testing.T) {
+       _, s := beam.NewPipelineWithRoot()
+       col := beam.Create(s, 1, 2, 3)
+       global := beam.Create(s, "ready")
+       fixed := beam.WindowInto(s, window.NewFixedWindows(time.Minute), 
beam.Create(s, "ready"))
+
+       // A global signal is fine in any position; the fixed-windowed one is 
not.
+       expectPanic(t, "wait.On: signal pcollection must be globally windowed 
when the input pcollection is (a global main window cannot be mapped to a 
non-global side-input window): index 1", func() {
+               wait.On(s, col, global, fixed)
+       })
+       // The other direction is allowed: a windowed main input may wait on a 
global signal.
+       windowed := beam.WindowInto(s, window.NewFixedWindows(time.Minute), col)
+       if out := wait.On(s, windowed, global); !out.IsValid() {
+               t.Errorf("wait.On(windowed main, global signal) returned an 
invalid PCollection")
+       }
+}
+
+// waitTestKVFn turns an int into a KV<int, string> entry so tests can build KV
+// PCollections from a package-level (non-closure) DoFn.
+func waitTestKVFn(v int) (int, string) {
+       return v, fmt.Sprintf("v%d", v)
+}
+
+// waitTestFormatKVFn renders a KV<int, string> as one string. passert rejects
+// composite element types, so KV outputs are compared in this projected form.
+func waitTestFormatKVFn(k int, v string) string {
+       return fmt.Sprintf("%d:%s", k, v)
+}
+
+// kvCol returns a KV<int, string> PCollection with one entry per value.
+func kvCol(s beam.Scope, values ...int) beam.PCollection {
+       return beam.ParDo(s, waitTestKVFn, beam.CreateList(s, values))
+}
+
+func TestOn_KVMainInput(t *testing.T) {
+       t.Run("PlainSignal", func(t *testing.T) {
+               ptest.BuildAndRun(t, func(s beam.Scope) {
+                       main := kvCol(s, 1, 2, 3)
+                       sig := beam.Create(s, "ready")
+                       out := wait.On(s, main, sig)
+                       passert.Equals(s, beam.ParDo(s, waitTestFormatKVFn, 
out), "1:v1", "2:v2", "3:v3")
+               })
+       })
+       t.Run("KVSignal", func(t *testing.T) {
+               ptest.BuildAndRun(t, func(s beam.Scope) {
+                       main := kvCol(s, 1, 2, 3)
+                       sig := kvCol(s, 10, 20)
+                       out := wait.On(s, main, sig)
+                       passert.Equals(s, beam.ParDo(s, waitTestFormatKVFn, 
out), "1:v1", "2:v2", "3:v3")
+               })
+       })
+}
+
+func TestOn_KVSignal(t *testing.T) {
+       ptest.BuildAndRun(t, func(s beam.Scope) {
+               main := beam.Create(s, 1, 2, 3)
+               sig := kvCol(s, 10, 20)
+               out := wait.On(s, main, sig)
+               passert.Equals(s, out, 1, 2, 3)
+       })
+}
+
+func TestOn_CoGBKPanics(t *testing.T) {
+       _, s := beam.NewPipelineWithRoot()
+       cogbk := beam.CoGroupByKey(s, kvCol(s, 1, 2, 3), kvCol(s, 4, 5, 6))
+       col := beam.Create(s, 1, 2, 3)
+       sig := beam.Create(s, "ready")
+
+       expectPanic(t, "wait.On: input pcollection must not be a CoGBK: ", 
func() {
+               wait.On(s, cogbk, sig)
+       })
+       expectPanic(t, "wait.On: signal pcollection must not be a CoGBK: index 
0: ", func() {
+               wait.On(s, col, cogbk)
+       })
+}
+
+// TestOn_PreservesCoder checks that the output retains the coder assigned to
+// the input through every chained stage.
+func TestOn_PreservesCoder(t *testing.T) {
+       t.Run("Plain", func(t *testing.T) {
+               _, s := beam.NewPipelineWithRoot()
+               col := beam.Create(s, 1, 2, 3)
+               replacement := beam.NewCoder(col.Type())
+               if err := col.SetCoder(replacement); err != nil {
+                       t.Fatalf("SetCoder failed: %v", err)
+               }
+               out := wait.On(s, col, beam.Create(s, "a"), beam.Create(s, 1.5))
+               if out.Coder() != replacement {
+                       t.Errorf("output coder %v is not the input's coder %v", 
out.Coder(), replacement)
+               }
+       })
+       t.Run("KV", func(t *testing.T) {
+               _, s := beam.NewPipelineWithRoot()
+               col := kvCol(s, 1, 2, 3)
+               replacement := beam.NewCoder(col.Type())
+               if err := col.SetCoder(replacement); err != nil {
+                       t.Fatalf("SetCoder failed: %v", err)
+               }
+               out := wait.On(s, col, beam.Create(s, "a"))
+               if out.Coder() != replacement {
+                       t.Errorf("output coder %v is not the input's coder %v", 
out.Coder(), replacement)
+               }
+       })
+}
diff --git a/sdks/go/test/integration/integration.go 
b/sdks/go/test/integration/integration.go
index 535a7fb8413..033a19bcdca 100644
--- a/sdks/go/test/integration/integration.go
+++ b/sdks/go/test/integration/integration.go
@@ -77,6 +77,8 @@ var directFilters = []string{
        // Triggers, Panes are not yet supported
        "TestTrigger.*",
        "TestPanes",
+       // TestStream-based wait tests run on Prism; bounded variants run here.
+       "TestWaitStream.*",
        // The direct runner does not support the TestStream primitive
        "TestTestStream.*",
        // (https://github.com/apache/beam/issues/21130): The direct runner 
does not support windowed side inputs
@@ -117,6 +119,8 @@ var portableFilters = []string{
        // The trigger and pane tests uses TestStream
        "TestTrigger.*",
        "TestPanes",
+       // The portable runner does not support the TestStream primitive; only 
bounded tests run here.
+       "TestWaitStream.*",
        // TODO(https://github.com/apache/beam/issues/21058): Python portable 
runner times out on Kafka reads.
        "TestKafkaIO.*",
        // TODO(BEAM-13215): GCP IOs currently do not work in non-Dataflow 
portable runners.
@@ -198,6 +202,10 @@ var flinkFilters = []string{
        "TestTestStreamSimple_InfinityDefault",
        "TestTestStreamToGBK",
        "TestTestStreamTimersEventTime",
+       // TODO(https://github.com/apache/beam/issues/31122): Flink's TestStream
+       // corrupts some length-prefixed and custom-coded values. These tests 
use a
+       // user-defined struct; sequencing variants also rely on process-local 
state.
+       "TestWaitStream.*",
 
        "TestTimers_EventTime_WithNoOutputTimestamp", // Encounter error: 
TimestampCombiner moved element from TIMESTAMP_MAX_VALUE to earlier time (end 
of global window) for window GlobalWindow
        "TestTimers_ProcessingTime.*",                // Flink doesn't support 
processing time timers.
@@ -216,6 +224,9 @@ var sparkFilters = []string{
        // The trigger and pane tests uses TestStream
        "TestTrigger.*",
        "TestPanes",
+       // Spark does not support TestStream or side inputs to executable 
stages.
+       "TestWaitStream.*",
+       "TestWaitBounded.*",
        // [BEAM-13921]: Spark doesn't support side inputs to executable stages
        "TestDebeziumIO_BasicRead",
        // TODO(BEAM-13215): GCP IOs currently do not work in non-Dataflow 
portable runners.
@@ -257,6 +268,8 @@ var dataflowFilters = []string{
        // The trigger and pane tests uses TestStream
        "TestTrigger.*",
        "TestPanes",
+       // The Dataflow runner does not support the TestStream primitive; only 
bounded tests run here.
+       "TestWaitStream.*",
        // There is no infrastructure for running KafkaIO tests with Dataflow.
        "TestKafkaIO.*",
        "TestSpannerIO.*",
diff --git a/sdks/go/test/integration/primitives/wait.go 
b/sdks/go/test/integration/primitives/wait.go
new file mode 100644
index 00000000000..3b0583cb17c
--- /dev/null
+++ b/sdks/go/test/integration/primitives/wait.go
@@ -0,0 +1,307 @@
+// Licensed to the Apache Software Foundation (ASF) under one or more
+// contributor license agreements.  See the NOTICE file distributed with
+// this work for additional information regarding copyright ownership.
+// The ASF licenses this file to You under the Apache License, Version 2.0
+// (the "License"); you may not use this file except in compliance with
+// the License.  You may obtain a copy of the License at
+//
+//    http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package primitives
+
+import (
+       "fmt"
+       "math"
+       "reflect"
+       "sync/atomic"
+       "time"
+
+       "github.com/apache/beam/sdks/v2/go/pkg/beam"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph/mtime"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph/window"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/register"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/testing/passert"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/testing/teststream"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/transforms/wait"
+)
+
+func init() {
+       // Package initialization also runs in separately launched SDK workers.
+       waitMaxMainTs.Store(math.MinInt64)
+
+       beam.RegisterType(reflect.TypeOf((*waitEvent)(nil)).Elem())
+       register.Function1x1(waitPartitionFn)
+       register.Function1x1(waitValueFn)
+       register.Function2x2(waitCheckSignalFn)
+       register.Function2x1(waitRecordMainFn)
+       register.Function2x0(waitTimestampFn)
+       register.Emitter2[beam.EventTime, waitEvent]()
+}
+
+// Prism permits one TestStream per pipeline, so Signal distinguishes its two
+// outputs.
+type waitEvent struct {
+       Signal bool
+       V      int
+}
+
+// waitStep adds an element or advances the watermark, in milliseconds.
+type waitStep struct {
+       elem *waitEvent // non-nil for an element step
+       ts   int64      // element timestamp, or the new watermark
+}
+
+func elemAt(ts int64, signal bool, v int) waitStep {
+       return waitStep{elem: &waitEvent{Signal: signal, V: v}, ts: ts}
+}
+
+func watermarkTo(ts int64) waitStep {
+       return waitStep{ts: ts}
+}
+
+func waitPartitionFn(e waitEvent) int {
+       if e.Signal {
+               return 1
+       }
+       return 0
+}
+
+func waitValueFn(e waitEvent) int {
+       return e.V
+}
+
+// waitMainAndSignal builds and splits the shared TestStream.
+func waitMainAndSignal(s beam.Scope, steps []waitStep) (main, signal 
beam.PCollection) {
+       con := teststream.NewConfig()
+       for _, st := range steps {
+               var err error
+               if st.elem != nil {
+                       err = con.AddElements(st.ts, *st.elem)
+               } else {
+                       err = con.AdvanceWatermark(st.ts)
+               }
+               if err != nil {
+                       panic(err)
+               }
+       }
+       if err := con.AdvanceWatermarkToInfinity(); err != nil {
+               panic(err)
+       }
+       parts := beam.Partition(s, 2, waitPartitionFn, teststream.Create(s, 
con))
+       return parts[0], parts[1]
+}
+
+func waitAssertValues(s beam.Scope, out beam.PCollection, want ...int) {
+       vals := beam.ParDo(s, waitValueFn, out)
+       vals = beam.WindowInto(s, window.NewGlobalWindows(), vals)
+       passert.EqualsList(s, vals, want)
+}
+
+// runWaitScenario verifies that waiting preserves the main input.
+func runWaitScenario(s beam.Scope, steps []waitStep, mainWfn, sigWfn 
*window.Fn, wantMain ...int) {
+       main, signal := waitMainAndSignal(s, steps)
+       main = beam.WindowInto(s, mainWfn, main)
+       signal = beam.WindowInto(s, sigWfn, signal)
+       waitAssertValues(s, wait.On(s, main, signal), wantMain...)
+}
+
+func waitThreeWindowSteps() []waitStep {
+       return []waitStep{
+               elemAt(1_000, false, 1), elemAt(2_000, true, 100),
+               watermarkTo(5_000),
+               elemAt(16_000, false, 2), elemAt(17_000, true, 200),
+               watermarkTo(20_000),
+               elemAt(31_000, false, 3), elemAt(31_000, true, 300),
+       }
+}
+
+// WaitStreamSameFixedWindows waits with main and signal in identical 15s 
windows.
+func WaitStreamSameFixedWindows(s beam.Scope) {
+       runWaitScenario(s, waitThreeWindowSteps(),
+               window.NewFixedWindows(15*time.Second),
+               window.NewFixedWindows(15*time.Second),
+               1, 2, 3)
+}
+
+// WaitStreamDifferentFixedWindows waits with a 15s main input on a 7s signal.
+//
+// This and WaitStreamSlidingSignal only check that nothing is lost or
+// duplicated when the WindowFns differ. Prism gates side inputs stage-wide by
+// watermark, not per mapped window, so the mapping cannot be observed here;
+// TestValidateWindowedSideInputs covers the mapping itself.
+func WaitStreamDifferentFixedWindows(s beam.Scope) {
+       runWaitScenario(s, waitThreeWindowSteps(),
+               window.NewFixedWindows(15*time.Second),
+               window.NewFixedWindows(7*time.Second),
+               1, 2, 3)
+}
+
+// WaitStreamSlidingSignal waits on 7s sliding windows emitted every second.
+func WaitStreamSlidingSignal(s beam.Scope) {
+       runWaitScenario(s, waitThreeWindowSteps(),
+               window.NewFixedWindows(15*time.Second),
+               window.NewSlidingWindows(1*time.Second, 7*time.Second),
+               1, 2, 3)
+}
+
+// WaitStreamSomeSignalWindowsEmpty verifies that empty signal windows unblock.
+func WaitStreamSomeSignalWindowsEmpty(s beam.Scope) {
+       steps := []waitStep{
+               elemAt(1_000, false, 1),
+               watermarkTo(10_000),
+               elemAt(11_000, false, 2), elemAt(12_000, true, 200),
+               watermarkTo(20_000),
+               elemAt(21_000, false, 3),
+       }
+       runWaitScenario(s, steps,
+               window.NewFixedWindows(10*time.Second),
+               window.NewFixedWindows(10*time.Second),
+               1, 2, 3)
+}
+
+// WaitBoundedGlobalWindow covers bounded global windows on supported runners.
+func WaitBoundedGlobalWindow(s beam.Scope) {
+       main := beam.Create(s, 1, 2, 3)
+       signal := beam.Create(s, "ready")
+       passert.Equals(s, wait.On(s, main, signal), 1, 2, 3)
+}
+
+func waitTimestampFn(e waitEvent, emit func(beam.EventTime, waitEvent)) {
+       emit(mtime.FromMilliseconds(int64(e.V)*1000), e)
+}
+
+// waitBoundedEvents timestamps each value at that many seconds.
+func waitBoundedEvents(s beam.Scope, vals []int) beam.PCollection {
+       events := make([]waitEvent, len(vals))
+       for i, v := range vals {
+               events[i] = waitEvent{V: v}
+       }
+       return beam.ParDo(s, waitTimestampFn, beam.CreateList(s, events))
+}
+
+// WaitBoundedFixedWindows runs Wait with differently sized fixed windows on
+// supported runners.
+func WaitBoundedFixedWindows(s beam.Scope) {
+       main := beam.WindowInto(s, window.NewFixedWindows(15*time.Second), 
waitBoundedEvents(s, []int{1, 16, 31}))
+       signal := beam.WindowInto(s, window.NewFixedWindows(7*time.Second), 
waitBoundedEvents(s, []int{2, 17, 31}))
+       waitAssertValues(s, wait.On(s, main, signal), 1, 16, 31)
+}
+
+// waitMaxMainTs coordinates sequencing assertions when both probe DoFns run in
+// the same SDK worker process, as they do in Prism's single worker 
environment.
+var waitMaxMainTs atomic.Int64
+
+// waitCheckSignalFn fails if an older signal arrives after main input was
+// released. It runs before wait.On's marker Combine.
+func waitCheckSignalFn(ts beam.EventTime, e waitEvent) (waitEvent, error) {
+       maxMain := waitMaxMainTs.Load()
+       if maxMain != math.MinInt64 && ts.Milliseconds() < maxMain {
+               return e, fmt.Errorf("signal element %+v at %v was processed 
after wait.On released main input up to %v",
+                       e, ts, mtime.FromMilliseconds(maxMain))
+       }
+       return e, nil
+}
+
+// waitRecordMainFn records the latest released main-input timestamp.
+func waitRecordMainFn(ts beam.EventTime, e waitEvent) waitEvent {
+       for {
+               cur := waitMaxMainTs.Load()
+               if ts.Milliseconds() <= cur || 
waitMaxMainTs.CompareAndSwap(cur, ts.Milliseconds()) {
+                       return e
+               }
+       }
+}
+
+// WaitStreamSequencing verifies that main@5s is not released before the older
+// signal@1s is processed. The allowed-lateness test below separately verifies
+// the Never trigger and expiration behavior.
+func WaitStreamSequencing(s beam.Scope) {
+       waitMaxMainTs.Store(math.MinInt64)
+
+       steps := []waitStep{
+               elemAt(5_000, false, 5),
+               elemAt(1_000, true, 1),
+               watermarkTo(10_000),
+               elemAt(15_000, false, 15),
+               elemAt(12_000, true, 12),
+               watermarkTo(20_000),
+       }
+       main, signal := waitMainAndSignal(s, steps)
+       wfn := window.NewFixedWindows(10 * time.Second)
+       main = beam.WindowInto(s, wfn, main)
+       signal = beam.WindowInto(s, wfn, beam.ParDo(s, waitCheckSignalFn, 
signal))
+
+       out := beam.ParDo(s, waitRecordMainFn, wait.On(s, main, signal))
+       waitAssertValues(s, out, 5, 15)
+}
+
+// waitLatenessSteps adds signal@3s after the watermark reaches 10s. It remains
+// valid until the 5s allowed lateness expires.
+func waitLatenessSteps() []waitStep {
+       return []waitStep{
+               elemAt(5_000, false, 5),
+               elemAt(2_000, true, 2),
+               watermarkTo(10_000),
+               elemAt(3_000, true, 3), // late, within the signal's 5s allowed 
lateness
+               watermarkTo(15_000),
+       }
+}
+
+var waitLatenessFixedWindows = window.NewFixedWindows(10 * time.Second)
+
+func waitAssertLatenessOutput(s beam.Scope, out beam.PCollection) {
+       waitAssertValues(s, beam.ParDo(s, waitRecordMainFn, out), 5)
+}
+
+// WaitStreamSequencingAllowedLateness verifies that [0,10s) remains blocked
+// until 15s when the signal has 5s allowed lateness. An early release causes
+// the late signal@3s probe to fail waitCheckSignalFn.
+func WaitStreamSequencingAllowedLateness(s beam.Scope) {
+       waitMaxMainTs.Store(math.MinInt64)
+
+       main, signal := waitMainAndSignal(s, waitLatenessSteps())
+       main = beam.WindowInto(s, waitLatenessFixedWindows, main)
+       signal = beam.WindowInto(s, waitLatenessFixedWindows, beam.ParDo(s, 
waitCheckSignalFn, signal),
+               beam.AllowedLateness(5*time.Second))
+
+       waitAssertLatenessOutput(s, wait.On(s, main, signal))
+}
+
+// waitStreamSequencingMultipleSignals verifies that every signal blocks,
+// regardless of argument order. The signals close at 10s and 15s.
+func waitStreamSequencingMultipleSignals(s beam.Scope, strictFirst bool) {
+       waitMaxMainTs.Store(math.MinInt64)
+
+       main, signal := waitMainAndSignal(s, waitLatenessSteps())
+       main = beam.WindowInto(s, waitLatenessFixedWindows, main)
+       checked := beam.ParDo(s, waitCheckSignalFn, signal)
+       strict := beam.WindowInto(s.Scope("strict"), waitLatenessFixedWindows, 
checked)
+       lenient := beam.WindowInto(s.Scope("lenient"), 
waitLatenessFixedWindows, checked,
+               beam.AllowedLateness(5*time.Second))
+
+       var out beam.PCollection
+       if strictFirst {
+               out = wait.On(s, main, strict, lenient)
+       } else {
+               out = wait.On(s, main, lenient, strict)
+       }
+       waitAssertLatenessOutput(s, out)
+}
+
+// WaitStreamSequencingMultipleSignals waits on a strict signal and then a
+// lenient one.
+func WaitStreamSequencingMultipleSignals(s beam.Scope) {
+       waitStreamSequencingMultipleSignals(s, true)
+}
+
+// WaitStreamSequencingMultipleSignalsLenientFirst waits on a lenient signal 
and
+// then a strict one.
+func WaitStreamSequencingMultipleSignalsLenientFirst(s beam.Scope) {
+       waitStreamSequencingMultipleSignals(s, false)
+}
diff --git a/sdks/go/test/integration/primitives/wait_test.go 
b/sdks/go/test/integration/primitives/wait_test.go
new file mode 100644
index 00000000000..e17e02c36c2
--- /dev/null
+++ b/sdks/go/test/integration/primitives/wait_test.go
@@ -0,0 +1,73 @@
+// Licensed to the Apache Software Foundation (ASF) under one or more
+// contributor license agreements.  See the NOTICE file distributed with
+// this work for additional information regarding copyright ownership.
+// The ASF licenses this file to You under the Apache License, Version 2.0
+// (the "License"); you may not use this file except in compliance with
+// the License.  You may obtain a copy of the License at
+//
+//    http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package primitives
+
+import (
+       "testing"
+
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/testing/ptest"
+       "github.com/apache/beam/sdks/v2/go/test/integration"
+)
+
+func TestWaitStreamSameFixedWindows(t *testing.T) {
+       integration.CheckFilters(t)
+       ptest.BuildAndRun(t, WaitStreamSameFixedWindows)
+}
+
+func TestWaitStreamDifferentFixedWindows(t *testing.T) {
+       integration.CheckFilters(t)
+       ptest.BuildAndRun(t, WaitStreamDifferentFixedWindows)
+}
+
+func TestWaitStreamSlidingSignal(t *testing.T) {
+       integration.CheckFilters(t)
+       ptest.BuildAndRun(t, WaitStreamSlidingSignal)
+}
+
+func TestWaitStreamSomeSignalWindowsEmpty(t *testing.T) {
+       integration.CheckFilters(t)
+       ptest.BuildAndRun(t, WaitStreamSomeSignalWindowsEmpty)
+}
+
+func TestWaitStreamSequencing(t *testing.T) {
+       integration.CheckFilters(t)
+       ptest.BuildAndRun(t, WaitStreamSequencing)
+}
+
+func TestWaitStreamSequencingAllowedLateness(t *testing.T) {
+       integration.CheckFilters(t)
+       ptest.BuildAndRun(t, WaitStreamSequencingAllowedLateness)
+}
+
+func TestWaitStreamSequencingMultipleSignals(t *testing.T) {
+       integration.CheckFilters(t)
+       ptest.BuildAndRun(t, WaitStreamSequencingMultipleSignals)
+}
+
+func TestWaitStreamSequencingMultipleSignalsLenientFirst(t *testing.T) {
+       integration.CheckFilters(t)
+       ptest.BuildAndRun(t, WaitStreamSequencingMultipleSignalsLenientFirst)
+}
+
+func TestWaitBoundedGlobalWindow(t *testing.T) {
+       integration.CheckFilters(t)
+       ptest.BuildAndRun(t, WaitBoundedGlobalWindow)
+}
+
+func TestWaitBoundedFixedWindows(t *testing.T) {
+       integration.CheckFilters(t)
+       ptest.BuildAndRun(t, WaitBoundedFixedWindows)
+}

Reply via email to