cool9850311 commented on code in PR #73420:
URL: https://github.com/apache/airflow/pull/73420#discussion_r4195213070


##########
go-sdk/pkg/execution/client.go:
##########
@@ -57,15 +71,63 @@ func translateApiError(err error, code string, sentinel 
error, key string) error
 // over the comm socket using msgpack-framed IPC instead of HTTP.
 type CoordinatorClient struct {
        comm *CoordinatorComm
+       // Bound at construction, not per call: the Execution API scopes the 
task
+       // state store to the caller's own task instance ("ti:self").
+       tiID string
 }
 
 var _ sdk.Client = (*CoordinatorClient)(nil)
 
 // NewCoordinatorClient creates a new client backed by the comm socket.
-func NewCoordinatorClient(comm *CoordinatorComm) *CoordinatorClient {
+func NewCoordinatorClient(comm *CoordinatorComm, tiID string) 
*CoordinatorClient {
        return &CoordinatorClient{
                comm: comm,
+               tiID: tiID,
+       }
+}
+
+// resolveDefaultExpiry returns nil ("never expires") for a retention of 0.
+// Only an absent env value falls back (the runtime was not launched by the
+// coordinator); a malformed one fails as it does in Python rather than 
silently
+// retaining keys for a different period.
+func resolveDefaultExpiry(now time.Time) (any, error) {
+       days := fallbackRetentionDays
+       if raw, ok := os.LookupEnv(defaultRetentionDaysEnv); ok {

Review Comment:
   Removed.



##########
go-sdk/pkg/execution/client.go:
##########
@@ -253,3 +315,183 @@ func (c *CoordinatorClient) PushXCom(
        _, err := c.comm.Communicate(ctx, msg)
        return err
 }
+
+// GetTaskState requests a task state value from the supervisor.
+func (c *CoordinatorClient) GetTaskState(ctx context.Context, key string) 
(any, error) {
+       resp, err := c.comm.Communicate(
+               ctx,
+               genmodels.GetTaskStateStore{TIID: c.tiID, Key: key},
+       )
+       if err != nil {
+               return nil, translateApiError(err, errCodeTaskStoreNotFound, 
sdk.TaskStateNotFound, key)
+       }
+
+       var result genmodels.TaskStateStoreResult
+       if err := decodeBody(resp, &result); err != nil {
+               return nil, fmt.Errorf("decoding task state result: %w", err)
+       }
+
+       return result.Value, nil
+}
+
+// UnmarshalJSONTaskState gets a task state value and unmarshals it into 
pointer.
+func (c *CoordinatorClient) UnmarshalJSONTaskState(

Review Comment:
   Done in 32512d0d27, close to your sketch: `actx.Client().TaskStateStore()` 
returns a `TaskStateStore` with `Get`, `UnmarshalJSONValue`, `Set(ctx, key, 
value, opts ...SetOption)`, `Delete` and `Clear`, and `sdk.WithRetention` 
replaces the second setter. Docs and the example bundle follow.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to