src/event/ngx_event_pipe.c - nginx-1.31.7 nginx/ @ 939334eff

Functions defined

Source code


  1. /*
  2. * Copyright (C) Igor Sysoev
  3. * Copyright (C) Nginx, Inc.
  4. */


  5. #include <ngx_config.h>
  6. #include <ngx_core.h>
  7. #include <ngx_event.h>
  8. #include <ngx_event_pipe.h>


  9. static ngx_int_t ngx_event_pipe_read_upstream(ngx_event_pipe_t *p);
  10. static ngx_int_t ngx_event_pipe_write_to_downstream(ngx_event_pipe_t *p);

  11. static ngx_int_t ngx_event_pipe_write_chain_to_temp_file(ngx_event_pipe_t *p);
  12. static ngx_inline void ngx_event_pipe_remove_shadow_links(ngx_buf_t *buf);
  13. static ngx_int_t ngx_event_pipe_drain_chains(ngx_event_pipe_t *p);


  14. ngx_int_t
  15. ‌ngx_event_pipe(ngx_event_pipe_t *p, ngx_int_t do_write)
  16. {
  17.     ngx_int_t     rc;
  18.     ngx_uint_t    flags;
  19.     ngx_event_t  *rev, *wev;

  20.     for ( ;; ) {
  21.         if (do_write) {
  22.             p->log->action = "sending to client";

  23.             rc = ngx_event_pipe_write_to_downstream(p);

  24.             if (rc == NGX_ABORT) {
  25.                 return NGX_ABORT;
  26.             }

  27.             if (rc == NGX_BUSY) {
  28.                 return NGX_OK;
  29.             }
  30.         }

  31.         p->read = 0;
  32.         p->upstream_blocked = 0;

  33.         p->log->action = "reading upstream";

  34.         if (ngx_event_pipe_read_upstream(p) == NGX_ABORT) {
  35.             return NGX_ABORT;
  36.         }

  37.         if (!p->read && !p->upstream_blocked) {
  38.             break;
  39.         }

  40.         do_write = 1;
  41.     }

  42.     if (p->upstream
  43.         && p->upstream->fd != (ngx_socket_t) -1)
  44.     {
  45.         rev = p->upstream->read;

  46.         flags = (rev->eof || rev->error) ? NGX_CLOSE_EVENT : 0;

  47.         if (ngx_handle_read_event(rev, flags) != NGX_OK) {
  48.             return NGX_ABORT;
  49.         }

  50.         if (!rev->delayed) {
  51.             if (rev->active && !rev->ready) {
  52.                 ngx_add_timer(rev, p->read_timeout);

  53.             } else if (rev->timer_set) {
  54.                 ngx_del_timer(rev);
  55.             }
  56.         }
  57.     }

  58.     if (p->downstream->fd != (ngx_socket_t) -1
  59.         && p->downstream->data == p->output_ctx)
  60.     {
  61.         wev = p->downstream->write;
  62.         if (ngx_handle_write_event(wev, p->send_lowat) != NGX_OK) {
  63.             return NGX_ABORT;
  64.         }

  65.         if (!wev->delayed) {
  66.             if (wev->active && !wev->ready) {
  67.                 ngx_add_timer(wev, p->send_timeout);

  68.             } else if (wev->timer_set) {
  69.                 ngx_del_timer(wev);
  70.             }
  71.         }
  72.     }

  73.     return NGX_OK;
  74. }


  75. static ngx_int_t
  76. ‌ngx_event_pipe_read_upstream(ngx_event_pipe_t *p)
  77. {
  78.     off_t         limit;
  79.     ssize_t       n, size;
  80.     ngx_int_t     rc;
  81.     ngx_buf_t    *b;
  82.     ngx_msec_t    delay;
  83.     ngx_chain_t  *chain, *cl, *ln;

  84.     if (p->upstream_eof || p->upstream_error || p->upstream_done
  85.         || p->upstream == NULL)
  86.     {
  87.         return NGX_OK;
  88.     }

  89. #if (NGX_THREADS)

  90.     if (p->aio) {
  91.         ngx_log_debug0(NGX_LOG_DEBUG_EVENT, p->log, 0,
  92.                        "pipe read upstream: aio");
  93.         return NGX_AGAIN;
  94.     }

  95.     if (p->writing) {
  96.         ngx_log_debug0(NGX_LOG_DEBUG_EVENT, p->log, 0,
  97.                        "pipe read upstream: writing");

  98.         rc = ngx_event_pipe_write_chain_to_temp_file(p);

  99.         if (rc != NGX_OK) {
  100.             return rc;
  101.         }
  102.     }

  103. #endif

  104.     ngx_log_debug1(NGX_LOG_DEBUG_EVENT, p->log, 0,
  105.                    "pipe read upstream: %d", p->upstream->read->ready);

  106.     for ( ;; ) {

  107.         if (p->upstream_eof || p->upstream_error || p->upstream_done) {
  108.             break;
  109.         }

  110.         if (p->preread_bufs == NULL && !p->upstream->read->ready) {
  111.             break;
  112.         }

  113.         if (p->preread_bufs) {

  114.             /* use the pre-read bufs if they exist */

  115.             chain = p->preread_bufs;
  116.             p->preread_bufs = NULL;
  117.             n = p->preread_size;

  118.             ngx_log_debug1(NGX_LOG_DEBUG_EVENT, p->log, 0,
  119.                            "pipe preread: %z", n);

  120.             if (n) {
  121.                 p->read = 1;
  122.             }

  123.         } else {

  124. #if (NGX_HAVE_KQUEUE)

  125.             /*
  126.              * kqueue notifies about the end of file or a pending error.
  127.              * This test allows not to allocate a buf on these conditions
  128.              * and not to call c->recv_chain().
  129.              */

  130.             if (p->upstream->read->available == 0
  131.                 && p->upstream->read->pending_eof
  132. #if (NGX_SSL)
  133.                 && !p->upstream->ssl
  134. #endif
  135.                 )
  136.             {
  137.                 p->upstream->read->ready = 0;
  138.                 p->upstream->read->eof = 1;
  139.                 p->upstream_eof = 1;
  140.                 p->read = 1;

  141.                 if (p->upstream->read->kq_errno) {
  142.                     p->upstream->read->error = 1;
  143.                     p->upstream_error = 1;
  144.                     p->upstream_eof = 0;

  145.                     ngx_log_error(NGX_LOG_ERR, p->log,
  146.                                   p->upstream->read->kq_errno,
  147.                                   "kevent() reported that upstream "
  148.                                   "closed connection");
  149.                 }

  150.                 break;
  151.             }
  152. #endif

  153.             if (p->limit_rate) {
  154.                 if (p->upstream->read->delayed) {
  155.                     break;
  156.                 }

  157.                 limit = (off_t) p->limit_rate * (ngx_time() - p->start_sec + 1)
  158.                         - p->read_length;

  159.                 if (limit <= 0) {
  160.                     p->upstream->read->delayed = 1;
  161.                     delay = (ngx_msec_t) (- limit * 1000 / p->limit_rate + 1);
  162.                     ngx_add_timer(p->upstream->read, delay);
  163.                     break;
  164.                 }

  165.             } else {
  166.                 limit = 0;
  167.             }

  168.             if (p->free_raw_bufs) {

  169.                 /* use the free bufs if they exist */

  170.                 chain = p->free_raw_bufs;
  171.                 if (p->single_buf) {
  172.                     p->free_raw_bufs = p->free_raw_bufs->next;
  173.                     chain->next = NULL;
  174.                 } else {
  175.                     p->free_raw_bufs = NULL;
  176.                 }

  177.             } else if (p->allocated < p->bufs.num) {

  178.                 /* allocate a new buf if it's still allowed */

  179.                 b = ngx_create_temp_buf(p->pool, p->bufs.size);
  180.                 if (b == NULL) {
  181.                     return NGX_ABORT;
  182.                 }

  183.                 p->allocated++;

  184.                 chain = ngx_alloc_chain_link(p->pool);
  185.                 if (chain == NULL) {
  186.                     return NGX_ABORT;
  187.                 }

  188.                 chain->buf = b;
  189.                 chain->next = NULL;

  190.             } else if (!p->cacheable
  191.                        && !p->downstream_error
  192.                        && p->downstream->data == p->output_ctx
  193.                        && p->downstream->write->ready
  194.                        && !p->downstream->write->delayed)
  195.             {
  196.                 /*
  197.                  * if the bufs are not needed to be saved in a cache and
  198.                  * a downstream is ready then write the bufs to a downstream
  199.                  */

  200.                 p->upstream_blocked = 1;

  201.                 ngx_log_debug0(NGX_LOG_DEBUG_EVENT, p->log, 0,
  202.                                "pipe downstream ready");

  203.                 break;

  204.             } else if (p->cacheable
  205.                        || p->temp_file->offset < p->max_temp_file_size)
  206.             {

  207.                 /*
  208.                  * if it is allowed, then save some bufs from p->in
  209.                  * to a temporary file, and add them to a p->out chain
  210.                  */

  211.                 rc = ngx_event_pipe_write_chain_to_temp_file(p);

  212.                 ngx_log_debug1(NGX_LOG_DEBUG_EVENT, p->log, 0,
  213.                                "pipe temp offset: %O", p->temp_file->offset);

  214.                 if (rc == NGX_BUSY) {
  215.                     break;
  216.                 }

  217.                 if (rc != NGX_OK) {
  218.                     return rc;
  219.                 }

  220.                 chain = p->free_raw_bufs;
  221.                 if (p->single_buf) {
  222.                     p->free_raw_bufs = p->free_raw_bufs->next;
  223.                     chain->next = NULL;
  224.                 } else {
  225.                     p->free_raw_bufs = NULL;
  226.                 }

  227.             } else {

  228.                 /* there are no bufs to read in */

  229.                 ngx_log_debug0(NGX_LOG_DEBUG_EVENT, p->log, 0,
  230.                                "no pipe bufs to read in");

  231.                 break;
  232.             }

  233.             n = p->upstream->recv_chain(p->upstream, chain, limit);

  234.             ngx_log_debug1(NGX_LOG_DEBUG_EVENT, p->log, 0,
  235.                            "pipe recv chain: %z", n);

  236.             if (p->free_raw_bufs) {
  237.                 chain->next = p->free_raw_bufs;
  238.             }
  239.             p->free_raw_bufs = chain;

  240.             if (n == NGX_ERROR) {
  241.                 p->upstream_error = 1;
  242.                 break;
  243.             }

  244.             if (n == NGX_AGAIN) {
  245.                 if (p->single_buf) {
  246.                     ngx_event_pipe_remove_shadow_links(chain->buf);
  247.                 }

  248.                 break;
  249.             }

  250.             p->read = 1;

  251.             if (n == 0) {
  252.                 p->upstream_eof = 1;
  253.                 break;
  254.             }
  255.         }

  256.         delay = p->limit_rate ? (ngx_msec_t) n * 1000 / p->limit_rate : 0;

  257.         p->read_length += n;
  258.         cl = chain;
  259.         p->free_raw_bufs = NULL;

  260.         while (cl && n > 0) {

  261.             ngx_event_pipe_remove_shadow_links(cl->buf);

  262.             size = cl->buf->end - cl->buf->last;

  263.             if (n >= size) {
  264.                 cl->buf->last = cl->buf->end;

  265.                 /* STUB */ cl->buf->num = p->num++;

  266.                 if (p->input_filter(p, cl->buf) == NGX_ERROR) {
  267.                     return NGX_ABORT;
  268.                 }

  269.                 n -= size;
  270.                 ln = cl;
  271.                 cl = cl->next;
  272.                 ngx_free_chain(p->pool, ln);

  273.             } else {
  274.                 cl->buf->last += n;
  275.                 n = 0;
  276.             }
  277.         }

  278.         if (cl) {
  279.             for (ln = cl; ln->next; ln = ln->next) { /* void */ }

  280.             ln->next = p->free_raw_bufs;
  281.             p->free_raw_bufs = cl;
  282.         }

  283.         if (delay > 0) {
  284.             p->upstream->read->delayed = 1;
  285.             ngx_add_timer(p->upstream->read, delay);
  286.             break;
  287.         }
  288.     }

  289. #if (NGX_DEBUG)

  290.     for (cl = p->busy; cl; cl = cl->next) {
  291.         ngx_log_debug8(NGX_LOG_DEBUG_EVENT, p->log, 0,
  292.                        "pipe buf busy s:%d t:%d f:%d "
  293.                        "%p, pos %p, size: %z "
  294.                        "file: %O, size: %O",
  295.                        (cl->buf->shadow ? 1 : 0),
  296.                        cl->buf->temporary, cl->buf->in_file,
  297.                        cl->buf->start, cl->buf->pos,
  298.                        cl->buf->last - cl->buf->pos,
  299.                        cl->buf->file_pos,
  300.                        cl->buf->file_last - cl->buf->file_pos);
  301.     }

  302.     for (cl = p->out; cl; cl = cl->next) {
  303.         ngx_log_debug8(NGX_LOG_DEBUG_EVENT, p->log, 0,
  304.                        "pipe buf out  s:%d t:%d f:%d "
  305.                        "%p, pos %p, size: %z "
  306.                        "file: %O, size: %O",
  307.                        (cl->buf->shadow ? 1 : 0),
  308.                        cl->buf->temporary, cl->buf->in_file,
  309.                        cl->buf->start, cl->buf->pos,
  310.                        cl->buf->last - cl->buf->pos,
  311.                        cl->buf->file_pos,
  312.                        cl->buf->file_last - cl->buf->file_pos);
  313.     }

  314.     for (cl = p->in; cl; cl = cl->next) {
  315.         ngx_log_debug8(NGX_LOG_DEBUG_EVENT, p->log, 0,
  316.                        "pipe buf in   s:%d t:%d f:%d "
  317.                        "%p, pos %p, size: %z "
  318.                        "file: %O, size: %O",
  319.                        (cl->buf->shadow ? 1 : 0),
  320.                        cl->buf->temporary, cl->buf->in_file,
  321.                        cl->buf->start, cl->buf->pos,
  322.                        cl->buf->last - cl->buf->pos,
  323.                        cl->buf->file_pos,
  324.                        cl->buf->file_last - cl->buf->file_pos);
  325.     }

  326.     for (cl = p->free_raw_bufs; cl; cl = cl->next) {
  327.         ngx_log_debug8(NGX_LOG_DEBUG_EVENT, p->log, 0,
  328.                        "pipe buf free s:%d t:%d f:%d "
  329.                        "%p, pos %p, size: %z "
  330.                        "file: %O, size: %O",
  331.                        (cl->buf->shadow ? 1 : 0),
  332.                        cl->buf->temporary, cl->buf->in_file,
  333.                        cl->buf->start, cl->buf->pos,
  334.                        cl->buf->last - cl->buf->pos,
  335.                        cl->buf->file_pos,
  336.                        cl->buf->file_last - cl->buf->file_pos);
  337.     }

  338.     ngx_log_debug1(NGX_LOG_DEBUG_EVENT, p->log, 0,
  339.                    "pipe length: %O", p->length);

  340. #endif

  341.     if (p->free_raw_bufs && p->length != -1) {
  342.         cl = p->free_raw_bufs;

  343.         if (cl->buf->last - cl->buf->pos >= p->length) {

  344.             p->free_raw_bufs = cl->next;

  345.             /* STUB */ cl->buf->num = p->num++;

  346.             if (p->input_filter(p, cl->buf) == NGX_ERROR) {
  347.                 return NGX_ABORT;
  348.             }

  349.             ngx_free_chain(p->pool, cl);
  350.         }
  351.     }

  352.     if (p->length == 0) {
  353.         p->upstream_done = 1;
  354.         p->read = 1;
  355.     }

  356.     if ((p->upstream_eof || p->upstream_error) && p->free_raw_bufs) {

  357.         /* STUB */ p->free_raw_bufs->buf->num = p->num++;

  358.         if (p->input_filter(p, p->free_raw_bufs->buf) == NGX_ERROR) {
  359.             return NGX_ABORT;
  360.         }

  361.         p->free_raw_bufs = p->free_raw_bufs->next;

  362.         if (p->free_bufs && p->buf_to_file == NULL) {
  363.             for (cl = p->free_raw_bufs; cl; cl = cl->next) {
  364.                 if (cl->buf->shadow == NULL) {
  365.                     ngx_pfree(p->pool, cl->buf->start);
  366.                 }
  367.             }
  368.         }
  369.     }

  370.     if (p->cacheable && (p->in || p->buf_to_file)) {

  371.         ngx_log_debug0(NGX_LOG_DEBUG_EVENT, p->log, 0,
  372.                        "pipe write chain");

  373.         rc = ngx_event_pipe_write_chain_to_temp_file(p);

  374.         if (rc != NGX_OK) {
  375.             return rc;
  376.         }
  377.     }

  378.     return NGX_OK;
  379. }


  380. static ngx_int_t
  381. ‌ngx_event_pipe_write_to_downstream(ngx_event_pipe_t *p)
  382. {
  383.     u_char            *prev;
  384.     size_t             bsize;
  385.     ngx_int_t          rc;
  386.     ngx_uint_t         flush, flushed, prev_last_shadow;
  387.     ngx_chain_t       *out, **ll, *cl;
  388.     ngx_connection_t  *downstream;

  389.     downstream = p->downstream;

  390.     ngx_log_debug1(NGX_LOG_DEBUG_EVENT, p->log, 0,
  391.                    "pipe write downstream: %d", downstream->write->ready);

  392. #if (NGX_THREADS)

  393.     if (p->writing) {
  394.         rc = ngx_event_pipe_write_chain_to_temp_file(p);

  395.         if (rc == NGX_ABORT) {
  396.             return NGX_ABORT;
  397.         }
  398.     }

  399. #endif

  400.     flushed = 0;

  401.     for ( ;; ) {
  402.         if (p->downstream_error) {
  403.             return ngx_event_pipe_drain_chains(p);
  404.         }

  405.         if (p->upstream_eof || p->upstream_error || p->upstream_done) {

  406.             /* pass the p->out and p->in chains to the output filter */

  407.             for (cl = p->busy; cl; cl = cl->next) {
  408.                 cl->buf->recycled = 0;
  409.             }

  410.             if (p->out) {
  411.                 ngx_log_debug0(NGX_LOG_DEBUG_EVENT, p->log, 0,
  412.                                "pipe write downstream flush out");

  413.                 for (cl = p->out; cl; cl = cl->next) {
  414.                     cl->buf->recycled = 0;
  415.                 }

  416.                 rc = p->output_filter(p->output_ctx, p->out);

  417.                 if (rc == NGX_ERROR) {
  418.                     p->downstream_error = 1;
  419.                     return ngx_event_pipe_drain_chains(p);
  420.                 }

  421.                 p->out = NULL;
  422.             }

  423.             if (p->writing) {
  424.                 break;
  425.             }

  426.             if (p->in) {
  427.                 ngx_log_debug0(NGX_LOG_DEBUG_EVENT, p->log, 0,
  428.                                "pipe write downstream flush in");

  429.                 for (cl = p->in; cl; cl = cl->next) {
  430.                     cl->buf->recycled = 0;
  431.                 }

  432.                 rc = p->output_filter(p->output_ctx, p->in);

  433.                 if (rc == NGX_ERROR) {
  434.                     p->downstream_error = 1;
  435.                     return ngx_event_pipe_drain_chains(p);
  436.                 }

  437.                 p->in = NULL;
  438.             }

  439.             ngx_log_debug0(NGX_LOG_DEBUG_EVENT, p->log, 0,
  440.                            "pipe write downstream done");

  441.             /* TODO: free unused bufs */

  442.             p->downstream_done = 1;
  443.             break;
  444.         }

  445.         if (downstream->data != p->output_ctx
  446.             || !downstream->write->ready
  447.             || downstream->write->delayed)
  448.         {
  449.             break;
  450.         }

  451.         /* bsize is the size of the busy recycled bufs */

  452.         prev = NULL;
  453.         bsize = 0;

  454.         for (cl = p->busy; cl; cl = cl->next) {

  455.             if (cl->buf->recycled) {
  456.                 if (prev == cl->buf->start) {
  457.                     continue;
  458.                 }

  459.                 bsize += cl->buf->end - cl->buf->start;
  460.                 prev = cl->buf->start;
  461.             }
  462.         }

  463.         ngx_log_debug1(NGX_LOG_DEBUG_EVENT, p->log, 0,
  464.                        "pipe write busy: %uz", bsize);

  465.         out = NULL;

  466.         if (bsize >= (size_t) p->busy_size) {
  467.             flush = 1;
  468.             goto flush;
  469.         }

  470.         flush = 0;
  471.         ll = NULL;
  472.         prev_last_shadow = 1;

  473.         for ( ;; ) {
  474.             if (p->out) {
  475.                 cl = p->out;

  476.                 if (cl->buf->recycled) {
  477.                     ngx_log_error(NGX_LOG_ALERT, p->log, 0,
  478.                                   "recycled buffer in pipe out chain");
  479.                 }

  480.                 p->out = p->out->next;

  481.             } else if (!p->cacheable && !p->writing && p->in) {
  482.                 cl = p->in;

  483.                 ngx_log_debug3(NGX_LOG_DEBUG_EVENT, p->log, 0,
  484.                                "pipe write buf ls:%d %p %z",
  485.                                cl->buf->last_shadow,
  486.                                cl->buf->pos,
  487.                                cl->buf->last - cl->buf->pos);

  488.                 if (cl->buf->recycled && prev_last_shadow) {
  489.                     if (bsize + cl->buf->end - cl->buf->start > p->busy_size) {
  490.                         flush = 1;
  491.                         break;
  492.                     }

  493.                     bsize += cl->buf->end - cl->buf->start;
  494.                 }

  495.                 prev_last_shadow = cl->buf->last_shadow;

  496.                 p->in = p->in->next;

  497.             } else {
  498.                 break;
  499.             }

  500.             cl->next = NULL;

  501.             if (out) {
  502.                 *ll = cl;
  503.             } else {
  504.                 out = cl;
  505.             }
  506.             ll = &cl->next;
  507.         }

  508.     flush:

  509.         ngx_log_debug2(NGX_LOG_DEBUG_EVENT, p->log, 0,
  510.                        "pipe write: out:%p, f:%ui", out, flush);

  511.         if (out == NULL) {

  512.             if (!flush) {
  513.                 break;
  514.             }

  515.             /* a workaround for AIO */
  516.             if (flushed++ > 10) {
  517.                 return NGX_BUSY;
  518.             }
  519.         }

  520.         rc = p->output_filter(p->output_ctx, out);

  521.         ngx_chain_update_chains(p->pool, &p->free, &p->busy, &out, p->tag);

  522.         if (rc == NGX_ERROR) {
  523.             p->downstream_error = 1;
  524.             return ngx_event_pipe_drain_chains(p);
  525.         }

  526.         for (cl = p->free; cl; cl = cl->next) {

  527.             if (cl->buf->temp_file) {
  528.                 if (p->cacheable || !p->cyclic_temp_file) {
  529.                     continue;
  530.                 }

  531.                 /* reset p->temp_offset if all bufs had been sent */

  532.                 if (cl->buf->file_last == p->temp_file->offset) {
  533.                     p->temp_file->offset = 0;
  534.                 }
  535.             }

  536.             /* TODO: free buf if p->free_bufs && upstream done */

  537.             /* add the free shadow raw buf to p->free_raw_bufs */

  538.             if (cl->buf->last_shadow) {
  539.                 if (ngx_event_pipe_add_free_buf(p, cl->buf->shadow) != NGX_OK) {
  540.                     return NGX_ABORT;
  541.                 }

  542.                 cl->buf->last_shadow = 0;
  543.             }

  544.             cl->buf->shadow = NULL;
  545.         }
  546.     }

  547.     return NGX_OK;
  548. }


  549. static ngx_int_t
  550. ‌ngx_event_pipe_write_chain_to_temp_file(ngx_event_pipe_t *p)
  551. {
  552.     ssize_t       size, bsize, n;
  553.     ngx_buf_t    *b;
  554.     ngx_uint_t    prev_last_shadow;
  555.     ngx_chain_t  *cl, *tl, *next, *out, **ll, **last_out, **last_free;

  556. #if (NGX_THREADS)

  557.     if (p->writing) {

  558.         if (p->aio) {
  559.             return NGX_AGAIN;
  560.         }

  561.         out = p->writing;
  562.         p->writing = NULL;

  563.         n = ngx_write_chain_to_temp_file(p->temp_file, NULL);

  564.         if (n == NGX_ERROR) {
  565.             return NGX_ABORT;
  566.         }

  567.         goto done;
  568.     }

  569. #endif

  570.     if (p->buf_to_file) {
  571.         out = ngx_alloc_chain_link(p->pool);
  572.         if (out == NULL) {
  573.             return NGX_ABORT;
  574.         }

  575.         out->buf = p->buf_to_file;
  576.         out->next = p->in;

  577.     } else {
  578.         out = p->in;
  579.     }

  580.     if (!p->cacheable) {

  581.         size = 0;
  582.         cl = out;
  583.         ll = NULL;
  584.         prev_last_shadow = 1;

  585.         ngx_log_debug1(NGX_LOG_DEBUG_EVENT, p->log, 0,
  586.                        "pipe offset: %O", p->temp_file->offset);

  587.         do {
  588.             bsize = cl->buf->last - cl->buf->pos;

  589.             ngx_log_debug4(NGX_LOG_DEBUG_EVENT, p->log, 0,
  590.                            "pipe buf ls:%d %p, pos %p, size: %z",
  591.                            cl->buf->last_shadow, cl->buf->start,
  592.                            cl->buf->pos, bsize);

  593.             if (prev_last_shadow
  594.                 && ((size + bsize > p->temp_file_write_size)
  595.                     || (p->temp_file->offset + size + bsize
  596.                         > p->max_temp_file_size)))
  597.             {
  598.                 break;
  599.             }

  600.             prev_last_shadow = cl->buf->last_shadow;

  601.             size += bsize;
  602.             ll = &cl->next;
  603.             cl = cl->next;

  604.         } while (cl);

  605.         ngx_log_debug1(NGX_LOG_DEBUG_EVENT, p->log, 0, "size: %z", size);

  606.         if (ll == NULL) {
  607.             return NGX_BUSY;
  608.         }

  609.         if (cl) {
  610.             p->in = cl;
  611.             *ll = NULL;

  612.         } else {
  613.             p->in = NULL;
  614.             p->last_in = &p->in;
  615.         }

  616.     } else {
  617.         p->in = NULL;
  618.         p->last_in = &p->in;
  619.     }

  620. #if (NGX_THREADS)
  621.     if (p->thread_handler) {
  622.         p->temp_file->thread_write = 1;
  623.         p->temp_file->file.thread_task = p->thread_task;
  624.         p->temp_file->file.thread_handler = p->thread_handler;
  625.         p->temp_file->file.thread_ctx = p->thread_ctx;
  626.     }
  627. #endif

  628.     n = ngx_write_chain_to_temp_file(p->temp_file, out);

  629.     if (n == NGX_ERROR) {
  630.         return NGX_ABORT;
  631.     }

  632. #if (NGX_THREADS)

  633.     if (n == NGX_AGAIN) {
  634.         p->writing = out;
  635.         p->thread_task = p->temp_file->file.thread_task;
  636.         return NGX_AGAIN;
  637.     }

  638. done:

  639. #endif

  640.     if (p->buf_to_file) {
  641.         p->temp_file->offset = p->buf_to_file->last - p->buf_to_file->pos;
  642.         n -= p->buf_to_file->last - p->buf_to_file->pos;
  643.         p->buf_to_file = NULL;
  644.         out = out->next;
  645.     }

  646.     if (n > 0) {
  647.         /* update previous buffer or add new buffer */

  648.         if (p->out) {
  649.             for (cl = p->out; cl->next; cl = cl->next) { /* void */ }

  650.             b = cl->buf;

  651.             if (b->file_last == p->temp_file->offset) {
  652.                 p->temp_file->offset += n;
  653.                 b->file_last = p->temp_file->offset;
  654.                 goto free;
  655.             }

  656.             last_out = &cl->next;

  657.         } else {
  658.             last_out = &p->out;
  659.         }

  660.         cl = ngx_chain_get_free_buf(p->pool, &p->free);
  661.         if (cl == NULL) {
  662.             return NGX_ABORT;
  663.         }

  664.         b = cl->buf;

  665.         ngx_memzero(b, sizeof(ngx_buf_t));

  666.         b->tag = p->tag;

  667.         b->file = &p->temp_file->file;
  668.         b->file_pos = p->temp_file->offset;
  669.         p->temp_file->offset += n;
  670.         b->file_last = p->temp_file->offset;

  671.         b->in_file = 1;
  672.         b->temp_file = 1;

  673.         *last_out = cl;
  674.     }

  675. free:

  676.     for (last_free = &p->free_raw_bufs;
  677.          *last_free != NULL;
  678.          last_free = &(*last_free)->next)
  679.     {
  680.         /* void */
  681.     }

  682.     for (cl = out; cl; cl = next) {
  683.         next = cl->next;

  684.         cl->next = p->free;
  685.         p->free = cl;

  686.         b = cl->buf;

  687.         if (b->last_shadow) {

  688.             tl = ngx_alloc_chain_link(p->pool);
  689.             if (tl == NULL) {
  690.                 return NGX_ABORT;
  691.             }

  692.             tl->buf = b->shadow;
  693.             tl->next = NULL;

  694.             *last_free = tl;
  695.             last_free = &tl->next;

  696.             b->shadow->pos = b->shadow->start;
  697.             b->shadow->last = b->shadow->start;

  698.             ngx_event_pipe_remove_shadow_links(b->shadow);
  699.         }
  700.     }

  701.     return NGX_OK;
  702. }


  703. /* the copy input filter */

  704. ngx_int_t
  705. ‌ngx_event_pipe_copy_input_filter(ngx_event_pipe_t *p, ngx_buf_t *buf)
  706. {
  707.     ngx_buf_t    *b;
  708.     ngx_chain_t  *cl;

  709.     if (buf->pos == buf->last) {
  710.         return NGX_OK;
  711.     }

  712.     if (p->upstream_done) {
  713.         ngx_log_debug0(NGX_LOG_DEBUG_EVENT, p->log, 0,
  714.                        "input data after close");
  715.         return NGX_OK;
  716.     }

  717.     if (p->length == 0) {
  718.         p->upstream_done = 1;

  719.         ngx_log_error(NGX_LOG_WARN, p->log, 0,
  720.                       "upstream sent more data than specified in "
  721.                       "\"Content-Length\" header");

  722.         return NGX_OK;
  723.     }

  724.     cl = ngx_chain_get_free_buf(p->pool, &p->free);
  725.     if (cl == NULL) {
  726.         return NGX_ERROR;
  727.     }

  728.     b = cl->buf;

  729.     ngx_memcpy(b, buf, sizeof(ngx_buf_t));
  730.     b->shadow = buf;
  731.     b->tag = p->tag;
  732.     b->last_shadow = 1;
  733.     b->recycled = 1;
  734.     buf->shadow = b;

  735.     ngx_log_debug1(NGX_LOG_DEBUG_EVENT, p->log, 0, "input buf #%d", b->num);

  736.     if (p->in) {
  737.         *p->last_in = cl;
  738.     } else {
  739.         p->in = cl;
  740.     }
  741.     p->last_in = &cl->next;

  742.     if (p->length == -1) {
  743.         return NGX_OK;
  744.     }

  745.     if (b->last - b->pos > p->length) {

  746.         ngx_log_error(NGX_LOG_WARN, p->log, 0,
  747.                       "upstream sent more data than specified in "
  748.                       "\"Content-Length\" header");

  749.         b->last = b->pos + p->length;
  750.         p->upstream_done = 1;

  751.         return NGX_OK;
  752.     }

  753.     p->length -= b->last - b->pos;

  754.     return NGX_OK;
  755. }


  756. static ngx_inline void
  757. ‌ngx_event_pipe_remove_shadow_links(ngx_buf_t *buf)
  758. {
  759.     ngx_buf_t  *b, *next;

  760.     b = buf->shadow;

  761.     if (b == NULL) {
  762.         return;
  763.     }

  764.     while (!b->last_shadow) {
  765.         next = b->shadow;

  766.         b->temporary = 0;
  767.         b->recycled = 0;

  768.         b->shadow = NULL;
  769.         b = next;
  770.     }

  771.     b->temporary = 0;
  772.     b->recycled = 0;
  773.     b->last_shadow = 0;

  774.     b->shadow = NULL;

  775.     buf->shadow = NULL;
  776. }


  777. ngx_int_t
  778. ‌ngx_event_pipe_add_free_buf(ngx_event_pipe_t *p, ngx_buf_t *b)
  779. {
  780.     ngx_chain_t  *cl;

  781.     cl = ngx_alloc_chain_link(p->pool);
  782.     if (cl == NULL) {
  783.         return NGX_ERROR;
  784.     }

  785.     if (p->buf_to_file && b->start == p->buf_to_file->start) {
  786.         b->pos = p->buf_to_file->last;
  787.         b->last = p->buf_to_file->last;

  788.     } else {
  789.         b->pos = b->start;
  790.         b->last = b->start;
  791.     }

  792.     b->shadow = NULL;

  793.     cl->buf = b;

  794.     if (p->free_raw_bufs == NULL) {
  795.         p->free_raw_bufs = cl;
  796.         cl->next = NULL;

  797.         return NGX_OK;
  798.     }

  799.     if (p->free_raw_bufs->buf->pos == p->free_raw_bufs->buf->last) {

  800.         /* add the free buf to the list start */

  801.         cl->next = p->free_raw_bufs;
  802.         p->free_raw_bufs = cl;

  803.         return NGX_OK;
  804.     }

  805.     /* the first free buf is partially filled, thus add the free buf after it */

  806.     cl->next = p->free_raw_bufs->next;
  807.     p->free_raw_bufs->next = cl;

  808.     return NGX_OK;
  809. }


  810. static ngx_int_t
  811. ‌ngx_event_pipe_drain_chains(ngx_event_pipe_t *p)
  812. {
  813.     ngx_chain_t  *cl, *tl;

  814.     for ( ;; ) {
  815.         if (p->out) {
  816.             cl = p->out;
  817.             p->out = NULL;

  818.         } else if (p->in) {
  819.             cl = p->in;
  820.             p->in = NULL;

  821.         } else {
  822.             return NGX_OK;
  823.         }

  824.         while (cl) {
  825.             if (cl->buf->last_shadow) {
  826.                 if (ngx_event_pipe_add_free_buf(p, cl->buf->shadow) != NGX_OK) {
  827.                     return NGX_ABORT;
  828.                 }

  829.                 cl->buf->last_shadow = 0;
  830.             }

  831.             cl->buf->shadow = NULL;
  832.             tl = cl->next;
  833.             cl->next = p->free;
  834.             p->free = cl;
  835.             cl = tl;
  836.         }
  837.     }
  838. }