DBIO

 view release on metacpan or  search on metacpan

lib/DBIO/Storage/Async.pm  view on Meta::CPAN


sub _query_async_pinned {
  croak 'Subclass must override _query_async_pinned($conn, $sql, $bind)';
}


sub _await_query_result {
  croak 'Subclass must override _await_query_result($conn, $sql, $bind)';
}


sub _await_conn_ready {
  croak 'Subclass must override _await_conn_ready($conn)';
}


sub _transform_sql {
  croak 'Subclass must override _transform_sql($sql)';
}


sub _post_insert_sql {
  croak 'Subclass must override _post_insert_sql';
}


sub txn_do_async {
  my ($self, $coderef, @args) = @_;

  my $fc = $self->future_class;

  return $self->pool->acquire_txn->then(sub {
    my $conn = shift;

    my $txn_ctx_class = $self->_txn_context_class;
    my $accessor = $self->_txn_conn_accessor;
    my $txn_ctx = $txn_ctx_class->new(
      storage   => $self,
      $accessor => $conn,
    );

    # BEGIN
    return $self->_query_async_pinned($conn, 'BEGIN', [])->then(sub {
      my $inner = eval { $coderef->($txn_ctx, @args) };
      if ($@) {
        my $error = $@;
        return $self->_query_async_pinned($conn, 'ROLLBACK', [])->then(sub {
          $self->pool->release($conn);
          $fc->fail($error);
        }, sub {
          my $rerr = shift;
          $self->pool->release($conn);
          $fc->fail($rerr);
        });
      }

      # karr #10: the coderef's Future is almost always the tail of a
      # ->then chain. Real Future holds a downstream sequence Future only
      # WEAKLY, so unless we keep a strong ref it gets GC'd the moment this
      # callback returns -- Future warns "lost a sequence Future",
      # COMMIT/ROLLBACK never fires, and the await loop busy-spins forever.
      # ->retain gives $chain_f a self-reference until it is ready, keeping
      # the whole chain alive. Not every future_class implements ->retain
      # (an immediately-resolved shim has no GC window), so guard the call.
      if (ref $inner && $inner->can('then')) {
        my $chain_f = $inner->then(sub {
          my @result = @_;
          return $self->_query_async_pinned($conn, 'COMMIT', [])->then(sub {
            $self->pool->release($conn);
            $fc->done(@result);
          }, sub {
            my $cerr = shift;
            $self->pool->release($conn);
            $fc->fail("COMMIT failed: $cerr");
          });
        }, sub {
          my $error = shift;
          return $self->_query_async_pinned($conn, 'ROLLBACK', [])->then(sub {
            $self->pool->release($conn);
            $fc->fail($error);
          }, sub {
            my $rerr = shift;
            $self->pool->release($conn);
            $fc->fail($rerr);
          });
        });
        $chain_f->retain if $chain_f->can('retain');
        return $chain_f;
      }
      else {
        # Coderef returned a plain value -- commit immediately
        return $self->_query_async_pinned($conn, 'COMMIT', [])->then(sub {
          $self->pool->release($conn);
          $fc->done($inner);
        }, sub {
          my $cerr = shift;
          $self->pool->release($conn);
          $fc->fail("COMMIT failed: $cerr");
        });
      }
    }, sub {
      my $err = shift;
      $self->pool->release($conn);
      $fc->fail("BEGIN failed: $err");
    });
  });
}


sub _txn_context_class { 'DBIO::Storage::Async::TransactionContext' }


sub _txn_conn_accessor { 'txn_conn' }


sub pipeline {
  my ($self, $coderef) = @_;

  my $fc = $self->future_class;

  return $self->pool->acquire->then(sub {



( run in 3.831 seconds using v1.01-cache-2.11-cpan-9169edd2b0e )